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..>`): 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" } // CLIAsked is a line mesh-cli was given on a node, and the account the node-engine says is asking. type CLIAsked struct { Line []string `json:"line"` Account string `json:"account"` UID uint32 `json:"uid"` } // CLIAnswer is what the controller answers: what the command printed, how it exited, whether it ran as the // controller's terminal and why, or why nothing ran. The node-engine hands it to mesh-cli as it is. 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. Each line on its own goroutine, bounded by HandlerTimeout, and its answer cut to one bus message. func (b OverNATS) ServeCLI(handle CLIHandler, logger *log.Logger) (func(), error) { done := make(chan struct{}) bind := func() (*nats.Subscription, error) { return b.Conn.QueueSubscribe(CLISubjects, "mesh-cli", func(msg *nats.Msg) { go b.answerCLI(msg, 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 } func (b OverNATS) answerCLI(msg *nats.Msg, handle CLIHandler, logger *log.Logger) { if msg.Reply == "" { return } var answer CLIAnswer node, ok := CLINode(msg.Subject) var asked CLIAsked switch { case !ok: answer = CLIRefusal("not a mesh-cli subject: " + msg.Subject) case json.Unmarshal(msg.Data, &asked) != nil || len(asked.Line) == 0: answer = CLIRefusal("the node-engine's request could not be read, so nothing ran") default: ctx, cancel := context.WithTimeout(context.Background(), HandlerTimeout) answer = handle(ctx, node, asked) cancel() } body := FitCLIAnswer(answer, b.Conn.MaxPayload()) if err := msg.Respond(body); err != nil && logger != nil { logger.Printf("mesh-cli on %s: the answer could not be sent: %v", node, err) } } // FitCLIAnswer is the answer as one bus message of at most limit bytes: what the command printed is cut, standard // output first, and the cut is said (ADR 0272 §5) — never silently short. func FitCLIAnswer(a CLIAnswer, limit int64) []byte { body, _ := json.Marshal(a) if limit <= 0 || int64(len(body)) <= limit { return body } 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] } body, _ = json.Marshal(a) return body }