- commandEnvironment takes the terminal as a bool instead of reading an empty verb as one; every non-terminal line names its verb and is stripped of MESH_CLI_TERMINAL, and a test with the mark set in the serving environment fails when that strip is taken out. - On the control-node the operator's account is the terminal only from a login session, as the node-engine reads it from the kernel's cgroup; the tool runner and the account's user units run as the operator too, and are ordinary calls. The request carries `session` (field-name tests on both sides). - A pull request the forge never announced is named with its number in the condition's headline (hq issue 347), from the stalled line's `number`, which mesh-delivery sends.
221 lines
8.1 KiB
Go
221 lines
8.1 KiB
Go
package link
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"log"
|
|
"strings"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
)
|
|
|
|
// mesh-cli asks the controller through the node-engine of the machine it runs on (novox/hq ADR 0272 §3). The engine
|
|
// reads the asking account from the kernel and asks on its own node's subject, which only that node's engine may
|
|
// publish (its grant is `mesh.control.<node>.>`): so the node is a fact the bus server enforces, and the account a
|
|
// fact the kernel gave. The engine's side holds the same field names (mesh-host internal/meshcli, a test on each
|
|
// side), as for its health statement.
|
|
|
|
// CLISubjects is every node's mesh-cli subject, which the serving controller answers.
|
|
const CLISubjects = "mesh.control.*.cli"
|
|
|
|
// CLISubject is one node's.
|
|
func CLISubject(node string) string { return "mesh.control." + node + ".cli" }
|
|
|
|
// CLISeat is what a mesh-cli line is recorded under in the calls record, beside the seats' verbs: its verb there is
|
|
// the node it was asked on.
|
|
const CLISeat = "mesh-cli"
|
|
|
|
// CLIAtOnce bounds the mesh-cli lines running at once, from every node together. A person types one at a time; one
|
|
// over the bound is answered busy at once, and nothing runs.
|
|
var CLIAtOnce = 8
|
|
|
|
// CLIAsked is a line mesh-cli was given on a node, and the account the node-engine says is asking — or, with
|
|
// Follow, the call a line runs as, asked again by the same account on the same node until it ends.
|
|
type CLIAsked struct {
|
|
Line []string `json:"line,omitempty"`
|
|
Account string `json:"account"`
|
|
UID uint32 `json:"uid"`
|
|
Follow string `json:"follow,omitempty"`
|
|
// Session is the login session the asking process runs in, as the node-engine read it from the kernel's cgroup
|
|
// (`session-<id>.scope` under the account's own slice), or empty: a service, a user unit. Only a login session is
|
|
// the terminal (novox/hq ADR 0272).
|
|
Session string `json:"session,omitempty"`
|
|
}
|
|
|
|
// CLIAnswer is what the controller answers, as the `result` of a call's answer: what the command printed, how it
|
|
// exited, whether it ran as the controller's terminal and why, or why nothing ran. A line still running is answered
|
|
// as every call is (`running`, `call`), and followed.
|
|
type CLIAnswer struct {
|
|
Stdout []byte `json:"stdout,omitempty"`
|
|
Stderr []byte `json:"stderr,omitempty"`
|
|
Exit int `json:"exit"`
|
|
Terminal bool `json:"terminal"`
|
|
Why string `json:"why,omitempty"`
|
|
Refused string `json:"refused,omitempty"`
|
|
Cut bool `json:"cut,omitempty"`
|
|
}
|
|
|
|
// CLIRefusal is an answer saying nothing ran, and why.
|
|
func CLIRefusal(why string) CLIAnswer { return CLIAnswer{Exit: 1, Refused: why} }
|
|
|
|
// CLINode is the node a mesh-cli subject names, or false for any other subject.
|
|
func CLINode(subject string) (string, bool) {
|
|
rest, ok := strings.CutPrefix(subject, "mesh.control.")
|
|
if !ok {
|
|
return "", false
|
|
}
|
|
node, tail, ok := strings.Cut(rest, ".")
|
|
if !ok || tail != "cli" || node == "" {
|
|
return "", false
|
|
}
|
|
return node, true
|
|
}
|
|
|
|
// CLIHandler answers one line from one node.
|
|
type CLIHandler func(ctx context.Context, node string, asked CLIAsked) CLIAnswer
|
|
|
|
// ServeCLI answers mesh-cli for every node until stopped: one queue group, so of two controllers during a handover
|
|
// one answers.
|
|
//
|
|
// **Every line is a call** (review of ADR 0272): run, recorded and answered as a seat's verb is (calls.go) — its
|
|
// caller answered within AnswerWithin that it is still running, with its call, and the line's answer kept for
|
|
// the asker to follow. The bus permits an answer for a minute only (broker.ResponseTTL); a line answered when it
|
|
// ends lost every answer after that, and mesh-cli said nothing ran of a line that had.
|
|
func (b OverNATS) ServeCLI(handle CLIHandler, logger *log.Logger) (func(), error) {
|
|
return b.serveCLI(Calls, CLIAtOnce, handle, logger)
|
|
}
|
|
|
|
func (b OverNATS) serveCLI(calls *CallLog, atOnce int, handle CLIHandler, logger *log.Logger) (func(), error) {
|
|
done := make(chan struct{})
|
|
slots := make(chan struct{}, atOnce)
|
|
bind := func() (*nats.Subscription, error) {
|
|
return b.Conn.QueueSubscribe(CLISubjects, "mesh-cli", func(msg *nats.Msg) {
|
|
go b.answerCLI(msg, calls, slots, handle, logger)
|
|
})
|
|
}
|
|
sub, err := bind()
|
|
if err != nil {
|
|
return nil, fmt.Errorf("serving %s: %w", CLISubjects, err)
|
|
}
|
|
go keepBound(sub, bind, CLISubjects, done, logger)
|
|
return func() {
|
|
close(done)
|
|
_ = sub.Unsubscribe()
|
|
}, nil
|
|
}
|
|
|
|
// cliEnvelope is an answer as every call's is: its result.
|
|
func cliEnvelope(a CLIAnswer) []byte {
|
|
body, _ := json.Marshal(map[string]any{"result": a})
|
|
return body
|
|
}
|
|
|
|
func (b OverNATS) answerCLI(msg *nats.Msg, calls *CallLog, slots chan struct{}, handle CLIHandler, logger *log.Logger) {
|
|
if msg.Reply == "" {
|
|
return
|
|
}
|
|
respond := func(body []byte) {
|
|
if err := msg.Respond(body); err != nil && logger != nil {
|
|
logger.Printf("mesh-cli on %s: the answer could not be sent: %v", msg.Subject, err)
|
|
}
|
|
}
|
|
node, ok := CLINode(msg.Subject)
|
|
var asked CLIAsked
|
|
switch {
|
|
case !ok:
|
|
respond(cliEnvelope(CLIRefusal("not a mesh-cli subject: " + msg.Subject)))
|
|
return
|
|
case json.Unmarshal(msg.Data, &asked) != nil:
|
|
respond(cliEnvelope(CLIRefusal("the node-engine's request could not be read, so nothing ran")))
|
|
return
|
|
case asked.Follow != "":
|
|
respond(calls.followCLI(asked.Follow, node, asked.Account))
|
|
return
|
|
case len(asked.Line) == 0:
|
|
respond(cliEnvelope(CLIRefusal("the request names no command, so nothing ran")))
|
|
return
|
|
}
|
|
select {
|
|
case slots <- struct{}{}:
|
|
defer func() { <-slots }()
|
|
default:
|
|
respond(cliEnvelope(CLIRefusal(fmt.Sprintf("busy: %d lines from mesh-cli are running already; ask again "+
|
|
"when one has ended. Nothing ran", cap(slots)))))
|
|
return
|
|
}
|
|
limit := b.Conn.MaxPayload()
|
|
calls.serveCallWithin(CLISeat, node, msg.Data, msg.Reply, func(ctx context.Context, _ json.RawMessage) (any, error) {
|
|
return fitCLI(handle(ctx, node, asked), limit-4096), nil
|
|
}, msg.Respond, limit, logger)
|
|
}
|
|
|
|
// followCLI answers the asker of a line what came of it: running still, its answer, or why that is not known. Only
|
|
// the account that asked it, on the node it was asked on: the answer is what the command printed, and `calls`
|
|
// keeps it from everyone else.
|
|
func (l *CallLog) followCLI(id, node, account string) []byte {
|
|
noSuch := func() []byte {
|
|
body, _ := json.Marshal(map[string]any{"error": fmt.Sprintf("no mesh-cli call %s was asked by %s on %s", id,
|
|
account, node)})
|
|
return body
|
|
}
|
|
c, inMemory := l.inMemory(id)
|
|
if !inMemory {
|
|
kept, found, err := l.Get(id)
|
|
if err != nil || !found {
|
|
return noSuch()
|
|
}
|
|
c = kept
|
|
}
|
|
var args CLIAsked
|
|
_ = json.Unmarshal(c.Args, &args)
|
|
if c.Seat != CLISeat || c.Verb != node || args.Account != account {
|
|
return noSuch()
|
|
}
|
|
switch {
|
|
case c.State == CallRunning:
|
|
return running(&c, AnswerWithin, false, l.Follow)
|
|
case c.State == CallAbandoned:
|
|
body, _ := json.Marshal(map[string]any{"error": fmt.Sprintf("call %s was running under a controller that "+
|
|
"stopped before it finished: it may have done part of what it was asked, and nothing will finish it", id)})
|
|
return body
|
|
case !inMemory:
|
|
body, _ := json.Marshal(map[string]any{"error": fmt.Sprintf("call %s finished (%s), and its answer was kept "+
|
|
"only in the memory of the controller that ran it, which has stopped; whether it took effect, the "+
|
|
"mesh says (status, the hand-act log)", id, c.State)})
|
|
return body
|
|
}
|
|
return c.Answer
|
|
}
|
|
|
|
// fitCLI is an answer whose streams fit in limit bytes once written: what the command printed is cut, standard
|
|
// output first, and the cut is said (ADR 0272 §5) — never silently short.
|
|
func fitCLI(a CLIAnswer, limit int64) CLIAnswer {
|
|
body, _ := json.Marshal(a)
|
|
if limit <= 0 || int64(len(body)) <= limit {
|
|
return a
|
|
}
|
|
a.Cut = true
|
|
// Room for the rest of the answer and JSON's base64 of the streams (4 bytes for every 3).
|
|
room := (limit - 4096) * 3 / 4
|
|
if room < 0 {
|
|
room = 0
|
|
}
|
|
if int64(len(a.Stderr)) > room/2 {
|
|
a.Stderr = a.Stderr[:room/2]
|
|
}
|
|
if left := room - int64(len(a.Stderr)); int64(len(a.Stdout)) > left {
|
|
if left < 0 {
|
|
left = 0
|
|
}
|
|
a.Stdout = a.Stdout[:left]
|
|
}
|
|
return a
|
|
}
|
|
|
|
// FitCLIAnswer is the answer as one bus message of at most limit bytes, written.
|
|
func FitCLIAnswer(a CLIAnswer, limit int64) []byte {
|
|
body, _ := json.Marshal(fitCLI(a, limit))
|
|
return body
|
|
}
|