diff --git a/cmd/mesh-controller/acts.go b/cmd/mesh-controller/acts.go index 6b1aa6da..9b06b742 100644 --- a/cmd/mesh-controller/acts.go +++ b/cmd/mesh-controller/acts.go @@ -51,6 +51,9 @@ func assignWith(ctx context.Context, open *stores, node string, opts assignOptio if len(modules) == 0 { return "", fmt.Errorf("assign %s names no module", node) } + if err := refusedMovingTheController(modules); err != nil { + return "", err + } // Held while it is recorded, so it cannot land between a converge's preview and its flip and // be taken without ever having been previewed (novox/hq ADR 0100). ctx, release, err := holdNodes(ctx, open, []string{node}) @@ -209,6 +212,9 @@ func seatDependenciesOnAssign(ctx context.Context, open *stores, node string, mo // (novox/hq ADR 0207) — the other side of refusing that module's assignment without one. Several // modules in one act are judged together, so a holder and its dependents come off in one command. func unassign(ctx context.Context, open *stores, node string, modules ...string) (string, error) { + if err := refusedMovingTheController(modules); err != nil { + return "", err + } if len(modules) == 0 { return "", fmt.Errorf("unassign %s names no module", node) } @@ -360,3 +366,17 @@ func issueOnAssign(ctx context.Context, open *stores, node, module string) strin } return fmt.Sprintf("its bus credential is issued and sealed to %s, and arrives with the push", node) } + +// refusedMovingTheController refuses assigning or unassigning the controller's module through a verb (review of +// novox/hq ADR 0272): the node it is assigned to is the control-node, and the control-node's operator account is +// the controller's terminal, so whoever moves the module chooses the terminal. At the terminal alone, as every +// change that says who may change what. +func refusedMovingTheController(modules []string) error { + verb, through := throughAVerb() + if !through || !slices.Contains(modules, controllerModule) { + return nil + } + return fmt.Errorf("%s is assigned and unassigned at the controller's terminal only (mesh-cli on the control-node), "+ + "never through a verb (this came through %q): the node it runs on is the control-node, whose operator is the "+ + "terminal (novox/hq ADR 0272). Nothing was assigned", controllerModule, verb) +} diff --git a/cmd/mesh-controller/meshcli.go b/cmd/mesh-controller/meshcli.go index 0e4cfa1d..f7b4c872 100644 --- a/cmd/mesh-controller/meshcli.go +++ b/cmd/mesh-controller/meshcli.go @@ -114,8 +114,24 @@ func answerMeshCLI(inv *inventory.Inventory) link.CLIHandler { } } -// runForMeshCLI runs a judged line and answers what it said. +// cliJournal says one line in the controller's journal (its standard output); a variable so a test can read it. +var cliJournal = func(line string) { fmt.Println(line) } + +// runForMeshCLI runs a judged line and answers what it said. Every line is said in the journal first: its call, +// who asked on which node, its command word only, and how it runs (review of ADR 0272). func runForMeshCLI(ctx context.Context, node string, asked link.CLIAsked, v cliVerdict) link.CLIAnswer { + how := "as an ordinary call" + switch { + case v.refused != "": + how = "refused: " + v.refused + case v.terminal: + how = "as the controller's terminal" + } + call := link.CallIDIn(ctx) + if call == "" { + call = "(no call)" + } + cliJournal(fmt.Sprintf("mesh-cli %s: %s on %s asked %q, %s", call, asked.Account, node, asked.Line[0], how)) if v.refused != "" { return link.CLIRefusal(v.refused) } @@ -123,14 +139,15 @@ func runForMeshCLI(ctx context.Context, node string, asked link.CLIAsked, v cliV return link.CLIAnswer{Exit: 1, Why: v.why, Refused: fmt.Sprintf("%s serves until stopped, and is not a "+ "command line mesh-cli runs. Nothing ran", asked.Line[0])} } - verb := "" + verb, line := "", asked.Line if !v.terminal { - if err := refusedAsTheGenericCommand(asked.Line); err != nil { + composed, err := ordinaryLine(asked.Line) + if err != nil { return link.CLIAnswer{Exit: 1, Why: v.why, Refused: err.Error()} } - verb = cliVerb + verb, line = cliVerb, composed } - cmd := selfCommand(ctx, asked.Line) + cmd := selfCommand(ctx, line) cmd.Env = commandEnvironment(fmt.Sprintf("%s through mesh-cli on %s", asked.Account, node), verb) // No standard input: a command that reads one gets nothing, and fails saying so (ADR 0272 §5). cmd.Stdin = nil @@ -155,3 +172,78 @@ func runForMeshCLI(ctx context.Context, node string, asked link.CLIAsked, v cliV } return answer } + +// settingsForms is what an ordinary `settings` line may say: the settings verb's own forms. +const settingsForms = "settings set [--replace] [--node ], settings clear " + + "[--node ], settings show [--history] [--node ], or settings preferences [] " + + "[--node ]" + +// ordinaryLine is the command line an ordinary call runs (ADR 0272 §4): a `settings` line composed exactly as the +// settings verb composes its own, so its refusals — the terminal-only keys among them — are that verb's (review of +// ADR 0272); any other line as the generic `command` verb takes it, with its refusals. +func ordinaryLine(argv []string) ([]string, error) { + if len(argv) == 0 || argv[0] != "settings" { + if err := refusedAsTheGenericCommand(argv); err != nil { + return nil, err + } + return argv, nil + } + args := map[string]any{} + var words []string + rest := argv[1:] + for i := 0; i < len(rest); i++ { + w := rest[i] + switch { + case w == "--replace" || w == "-replace": + args["replace"] = "true" + case w == "--history" || w == "-history": + args["history"] = "true" + case w == "--node" || w == "-node": + if i+1 >= len(rest) { + return nil, fmt.Errorf("--node names no node; settings takes %s. Nothing ran", settingsForms) + } + i++ + args["node"] = rest[i] + case strings.HasPrefix(w, "--node=") || strings.HasPrefix(w, "-node="): + _, args["node"], _ = strings.Cut(w, "=") + case strings.HasPrefix(w, "-"): + return nil, fmt.Errorf("settings takes no %s through mesh-cli outside the terminal; it takes %s. "+ + "Nothing ran", w, settingsForms) + default: + words = append(words, w) + } + } + wrong := fmt.Errorf("through mesh-cli outside the terminal, settings takes the settings verb's forms: %s. "+ + "Nothing ran", settingsForms) + if len(words) == 0 { + return nil, wrong + } + switch words[0] { + case "set": + if len(words) != 3 { + return nil, wrong + } + args["module"], args["values"] = words[1], words[2] + case "clear": + if len(words) != 2 { + return nil, wrong + } + args["module"], args["clear"] = words[1], "true" + case "show": + if len(words) != 2 { + return nil, wrong + } + args["module"] = words[1] + case "preferences": + if len(words) > 2 { + return nil, wrong + } + args["list"] = "preferences" + if len(words) == 2 { + args["module"] = words[1] + } + default: + return nil, wrong + } + return argvFor("settings", args) +} diff --git a/cmd/mesh-controller/meshcli_test.go b/cmd/mesh-controller/meshcli_test.go index e7b7eed2..0e9b82b7 100644 --- a/cmd/mesh-controller/meshcli_test.go +++ b/cmd/mesh-controller/meshcli_test.go @@ -84,15 +84,13 @@ func TestAnOrdinaryCallMeetsTheCommandVerbsRefusals(t *testing.T) { t.Setenv(echoEnvironment, "1") ctx := context.Background() ordinary := cliVerdict{why: "not the terminal"} - a := runForMeshCLI(ctx, "laptop", asked("operator", 1000, "settings", "set", "claude-code", "{}"), ordinary) - if a.Refused == "" || len(a.Stdout) != 0 || a.Exit != 1 { - t.Fatalf("settings set ran as an ordinary call: %+v", a) + a := runForMeshCLI(ctx, "laptop", asked("operator", 1000, "cleanup", "delete", "x"), ordinary) + if a.Refused == "" || len(a.Stdout) != 0 || a.Exit != 1 || a.Why != "not the terminal" { + t.Fatalf("a repair without --why ran as an ordinary call: %+v", a) } - if !strings.Contains(a.Refused, "controller's terminal") || a.Why != "not the terminal" { - t.Fatalf("the refusal does not say why: %+v", a) - } - if err := refusedAsTheGenericCommand([]string{"settings", "set", "x", "{}"}); err == nil { - t.Fatal("the command verb no longer refuses settings set, and mesh-cli's ordinary call relies on it") + a = runForMeshCLI(ctx, "laptop", asked("operator", 1000, "settings", "set", "claude-code", "{}"), ordinary) + if a.Refused != "" || !strings.Contains(string(a.Stdout), `verb="mesh-cli"`) { + t.Fatalf("an ordinary settings set did not run through the settings verb's path with MESH_VERB set: %+v", a) } for _, server := range []string{"serve", "api", "board"} { a := runForMeshCLI(ctx, "control", asked("operator", 1000, server), cliVerdict{terminal: true}) @@ -105,3 +103,77 @@ func TestAnOrdinaryCallMeetsTheCommandVerbsRefusals(t *testing.T) { t.Fatalf("a refused line ran: %+v", a) } } + +// **Which node is the terminal does not follow a verb** (review of ADR 0272): the controller's module is assigned +// and unassigned at the terminal alone, so no caller of `assign` can move the terminal to a node of its choosing. +func TestTheControllersModuleIsMovedAtTheTerminalAlone(t *testing.T) { + t.Setenv("MESH_VERB", "assign") + if err := refusedMovingTheController([]string{"zsh", "mesh-controller"}); err == nil || + !strings.Contains(err.Error(), "terminal") { + t.Fatalf("assigning the controller through a verb was not refused: %v", err) + } + if err := refusedMovingTheController([]string{"zsh"}); err != nil { + t.Fatalf("another module was refused: %v", err) + } + t.Setenv("MESH_VERB", "") + if err := refusedMovingTheController([]string{"mesh-controller"}); err != nil { + t.Fatalf("the terminal was refused: %v", err) + } +} + +// An ordinary `settings set|clear` goes down the settings verb's own path: composed as that verb composes it, and +// run with MESH_VERB set, so its refusals — the terminal-only keys among them — are the settings command's own, +// not a blanket refusal of the generic command (review of ADR 0272). +func TestAnOrdinarySettingsLineTakesTheSettingsVerbsPath(t *testing.T) { + cases := []struct { + line []string + want string + err string + }{ + {[]string{"settings", "set", "zsh", `{"execute":"withhold"}`, "--node", "laptop"}, `settings set zsh {"execute":"withhold"} --node laptop`, ""}, + {[]string{"settings", "set", "zsh", "{}", "--replace"}, "settings set zsh {} --replace", ""}, + {[]string{"settings", "clear", "zsh", "--node", "laptop"}, "settings clear zsh --node laptop", ""}, + {[]string{"settings", "show", "zsh", "--history"}, "settings show zsh --history", ""}, + {[]string{"settings", "preferences"}, "settings preferences", ""}, + {[]string{"settings", "set", "zsh"}, "", "settings"}, + {[]string{"settings", "set", "zsh", "{}", "--sideways"}, "", "--sideways"}, + } + for _, c := range cases { + argv, err := ordinaryLine(c.line) + if c.err != "" { + if err == nil || !strings.Contains(err.Error(), c.err) { + t.Errorf("%q: refused with %v, want %q", c.line, err, c.err) + } + continue + } + if err != nil || strings.Join(argv, " ") != c.want { + t.Errorf("%q: composed %q (%v), want %q", c.line, argv, err, c.want) + } + } + // Anything else still meets the generic command verb's refusals. + if _, err := ordinaryLine([]string{"retire", "delete", "x"}); err != nil && !strings.Contains(err.Error(), "why") { + t.Errorf("a repair was refused for another reason: %v", err) + } +} + +// Every line is said in the controller's journal, with its call, who asked where and how it ran — its command word +// only, never the rest of the line (review of ADR 0272). +func TestEveryMeshCLILineIsSaidInTheJournal(t *testing.T) { + t.Setenv(echoEnvironment, "1") + var said []string + was := cliJournal + cliJournal = func(line string) { said = append(said, line) } + t.Cleanup(func() { cliJournal = was }) + ctx := link.WithCallID(context.Background(), "call-1") + runForMeshCLI(ctx, "control", asked("operator", 1000, "settings", "set", "x", `{"password":"s3cret"}`), + cliVerdict{terminal: true, why: "the terminal"}) + runForMeshCLI(ctx, "control", asked("agent", 1001, "status"), cliVerdict{refused: "agent is not answered"}) + all := strings.Join(said, "\n") + if len(said) != 2 || !strings.Contains(all, "call-1") || !strings.Contains(all, "operator on control") || + !strings.Contains(all, "as the controller's terminal") || !strings.Contains(all, "refused") { + t.Fatalf("the journal said %q", said) + } + if strings.Contains(all, "s3cret") { + t.Fatalf("the journal carries the line's values: %q", said) + } +} diff --git a/cmd/mesh-controller/seatverbs.go b/cmd/mesh-controller/seatverbs.go index b0ba7b01..72a84ee1 100644 --- a/cmd/mesh-controller/seatverbs.go +++ b/cmd/mesh-controller/seatverbs.go @@ -1184,7 +1184,7 @@ func answersFirst(argv []string) bool { // the reason said beside it. func callsAnswer(log *link.CallLog, id string) (any, error) { if id != "" { - c, ok, err := log.Get(id) + c, ok, err := log.Shown(id) if err != nil { return nil, fmt.Errorf("call %s is not in this controller's memory, and the calls kept on the "+ "bus could not be read: %w", id, err) diff --git a/internal/broker/writers.go b/internal/broker/writers.go index 70cb099b..a775d4ce 100644 --- a/internal/broker/writers.go +++ b/internal/broker/writers.go @@ -87,8 +87,10 @@ var WritersTable = []WriterRow{ {State: "a machine's applied state and its report", Writer: "the node-engine's apply queue", KeptIn: "the machine; the report on the bus", Others: "the reconcile and a delivery enqueue, never apply", // And its health statement between reports (novox/hq ADR 0240): the same writer stating the same - // machine, inside the grant it already had (`mesh.control..>`). - Subjects: []string{"mesh.control.*.report", "mesh.control.*.health"}, Writes: ownMachine}, + // machine, inside the grant it already had (`mesh.control..>`). And a line mesh-cli was given + // on the machine (novox/hq ADR 0272): the node is what the controller judges the terminal by, so that + // subject is this machine's engine's alone. + Subjects: []string{"mesh.control.*.report", "mesh.control.*.health", "mesh.control.*.cli"}, Writes: ownMachine}, {State: "the controller lease", Writer: "the controller instance holding it", KeptIn: "key-value " + LeaseBucket, Others: "a candidate waits", Subjects: kvOf(LeaseBucket), Writes: isController}, {State: "plans and their tiers", Writer: "controller (lease holder), compare-and-set on the plan's revision", diff --git a/internal/broker/writers_test.go b/internal/broker/writers_test.go index 72aba5e4..11eeee7f 100644 --- a/internal/broker/writers_test.go +++ b/internal/broker/writers_test.go @@ -86,6 +86,10 @@ func TestASecondWriterIsRefusedAtComposition(t *testing.T) { []string{"mesh.control.>"}, "a machine's applied state and its report"}, {"a machine publishing another's report", Principal{Kind: KindNode, Node: "one"}, []string{"mesh.control.two.report"}, "a machine's applied state and its report"}, + {"a machine asking the controller as another (mesh-cli, ADR 0272)", Principal{Kind: KindNode, Node: "one"}, + []string{"mesh.control.two.cli"}, "a machine's applied state and its report"}, + {"a module asking the controller as a machine (mesh-cli, ADR 0272)", Principal{Kind: KindModule, Module: "shop"}, + []string{"mesh.control.one.cli"}, "a machine's applied state and its report"}, {"a machine publishing every machine's", Principal{Kind: KindNode, Node: "one"}, []string{"mesh.control.*.>"}, "a machine's applied state and its report"}, {"a module writing the lease", Principal{Kind: KindModule, Module: "shop"}, @@ -115,6 +119,7 @@ func TestASecondWriterIsRefusedAtComposition(t *testing.T) { publish []string }{ {Principal{Kind: KindNode, Node: "one"}, []string{"mesh.control.one.>"}}, + {Principal{Kind: KindNode, Node: "one"}, []string{"mesh.control.one.cli"}}, {Principal{Kind: KindController}, []string{"mesh.node.>", "$JS.API.>", "$KV.mesh-controller_lease.>"}}, {Principal{Kind: KindModule, Module: "gitea"}, []string{"mesh.mod.gitea.event.pull.merged"}}, {Principal{Kind: KindModule, Module: "shop"}, []string{"$JS.API.CONSUMER.CREATE.KV_shop_carts.>"}}, diff --git a/internal/link/calls.go b/internal/link/calls.go index daaaf6e1..29bbdc5a 100644 --- a/internal/link/calls.go +++ b/internal/link/calls.go @@ -324,6 +324,15 @@ func (l *CallLog) finish(c *Call, answer []byte, failed, answeredAlready bool) { return } onTheBus := *c + if c.Seat == CLISeat { + // **A mesh-cli line's answer is its asker's alone** (ADR 0272): what a command at the controller's terminal + // printed — a token, a secret's reference — and `calls` answers anyone who may call the seat. Kept in this + // process's memory for the asker to follow, and never on the bus. + onTheBus.Answer, _ = json.Marshal(map[string]any{"not kept": "a mesh-cli line's answer, its asker's alone: " + + "kept in the memory of the controller that ran it, for the asker to follow"}) + l.keep(onTheBus) + return + } if len(onTheBus.Answer) > keptOnTheBusAtMost { // A record on the bus is one message too, and one larger than the bus carries is refused whole — // the call's state with it (novox/hq issue 314). The answer stays in this process's memory, and @@ -391,6 +400,29 @@ func (l *CallLog) Recent() ([]Call, error) { return out, nil } +// inMemory is one call this process served and still holds, whole. +func (l *CallLog) inMemory(id string) (Call, bool) { + l.mu.Lock() + defer l.mu.Unlock() + for _, c := range l.calls { + if c.ID == id { + return *c, true + } + } + return Call{}, false +} + +// Shown is one kept call as `calls` shows it: whole, but for a mesh-cli line's answer, which is its asker's alone +// and followed by it (ADR 0272) — `calls` answers anyone who may call the seat. +func (l *CallLog) Shown(id string) (Call, bool, error) { + c, found, err := l.Get(id) + if found && c.Seat == CLISeat { + c.Answer, _ = json.Marshal(map[string]any{"not shown": "a mesh-cli line's answer is its asker's alone: " + + "mesh-cli follows it"}) + } + return c, found, err +} + // Get is one kept call: from memory, or from the bus when this process did not serve it. func (l *CallLog) Get(id string) (Call, bool, error) { l.mu.Lock() @@ -431,6 +463,11 @@ func kept(args json.RawMessage) json.RawMessage { switch { case k == "values" || k == "secret": out[k] = "(given, not kept)" + case k == "line": + // A mesh-cli line (ADR 0272): its command word, never the rest, which may carry settings. + if words, ok := v.([]any); ok && len(words) > 0 { + out[k] = []any{words[0], "(the rest given, not kept)"} + } case isString && len(s) <= 120: out[k] = s case isString: @@ -513,6 +550,19 @@ var watched struct { conns map[*nats.Conn]bool } +type callIDKey struct{} + +// WithCallID is ctx carrying the call it serves; CallIDIn reads it back, empty outside one. +func WithCallID(ctx context.Context, id string) context.Context { + return context.WithValue(ctx, callIDKey{}, id) +} + +// CallIDIn is the call ctx serves, or empty. +func CallIDIn(ctx context.Context) string { + id, _ := ctx.Value(callIDKey{}).(string) + return id +} + // answerNow is how a handler asks for its caller to be answered before it goes on (Acknowledge). type answerNowKey struct{} @@ -565,6 +615,8 @@ func (l *CallLog) serveCallWithin(seat, verb string, args json.RawMessage, reply c := l.begin(seat, verb, kept(args), reply) // Who asked travels with the call, so an act it does by hand says so (novox/hq to-be 45 §7). ctx = context.WithValue(ctx, callerKey{}, c.Caller) + // And which call it is, for what the handler says of it in the journal. + ctx = WithCallID(ctx, c.ID) acknowledged := make(chan struct{}) var once sync.Once ctx = context.WithValue(ctx, answerNowKey{}, func() { once.Do(func() { close(acknowledged) }) }) diff --git a/internal/link/meshcli.go b/internal/link/meshcli.go index a834f225..27dcc6f3 100644 --- a/internal/link/meshcli.go +++ b/internal/link/meshcli.go @@ -22,15 +22,26 @@ 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. +// 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"` + Line []string `json:"line,omitempty"` Account string `json:"account"` UID uint32 `json:"uid"` + Follow string `json:"follow,omitempty"` } -// 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. +// 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"` @@ -61,12 +72,22 @@ func CLINode(subject string) (string, bool) { 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. +// 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, handle, logger) + go b.answerCLI(msg, calls, slots, handle, logger) }) } sub, err := bind() @@ -80,35 +101,95 @@ func (b OverNATS) ServeCLI(handle CLIHandler, logger *log.Logger) (func(), error }, nil } -func (b OverNATS) answerCLI(msg *nats.Msg, handle CLIHandler, logger *log.Logger) { +// 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 } - var answer CLIAnswer + 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: - 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") + 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: - 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) + 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) } -// FitCLIAnswer is the answer as one bus message of at most limit bytes: what the command printed is cut, standard +// 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 FitCLIAnswer(a CLIAnswer, limit int64) []byte { +func fitCLI(a CLIAnswer, limit int64) CLIAnswer { body, _ := json.Marshal(a) if limit <= 0 || int64(len(body)) <= limit { - return body + return a } a.Cut = true // Room for the rest of the answer and JSON's base64 of the streams (4 bytes for every 3). @@ -125,6 +206,11 @@ func FitCLIAnswer(a CLIAnswer, limit int64) []byte { } a.Stdout = a.Stdout[:left] } - body, _ = json.Marshal(a) + 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 } diff --git a/internal/link/meshcli_test.go b/internal/link/meshcli_test.go index b56f20f2..d255be6d 100644 --- a/internal/link/meshcli_test.go +++ b/internal/link/meshcli_test.go @@ -1,10 +1,16 @@ package link import ( + "context" "encoding/json" "sort" "strings" "testing" + "time" + + "github.com/nats-io/nats.go" + + "github.com/novox/mesh-controller/internal/testbus" ) // The node-engine's request and the answer it hands mesh-cli hold these field names; the engine's side holds the @@ -19,6 +25,18 @@ func TestTheMeshCLIRequestAndAnswerKeepTheirFieldNames(t *testing.T) { if got := keysIn(t, body); got != "cut exit refused stderr stdout terminal why" { t.Fatalf("the answer's fields are %q", got) } + body, _ = json.Marshal(CLIAsked{Follow: "call-1", Account: "a", UID: 1}) + if got := keysIn(t, body); got != "account follow uid" { + t.Fatalf("a follow's fields are %q", got) + } + // What a line still running answers, in the envelope every call's answer is in: `result`, then these. + var env struct { + Result json.RawMessage `json:"result"` + } + _ = json.Unmarshal(running(&Call{ID: "call-1", Seat: CLISeat, Verb: "a"}, AnswerWithin, false, ""), &env) + if got := keysIn(t, env.Result); got != "call output running started" { + t.Fatalf("a running answer's fields are %q", got) + } } func keysIn(t *testing.T, body []byte) string { @@ -65,3 +83,166 @@ func TestALargeMeshCLIAnswerIsCutAndSaysSo(t *testing.T) { t.Fatal("an answer that fits was said to be cut") } } + +// cliBus serves mesh-cli on a bus of the test's own with a call log of its own, and returns a connection to ask on. +func cliBus(t *testing.T, l *CallLog, atOnce int, handle CLIHandler) *nats.Conn { + t.Helper() + conn, err := nats.Connect(testbus.URL(t)) + if err != nil { + t.Fatal(err) + } + t.Cleanup(conn.Close) + stop, err := OverNATS{Conn: conn}.serveCLI(l, atOnce, handle, nil) + if err != nil { + t.Fatal(err) + } + t.Cleanup(stop) + return conn +} + +// askCLI asks one request on a node's subject and reads the envelope. +func askCLI(t *testing.T, conn *nats.Conn, node string, asked CLIAsked) (CLIAnswer, map[string]any, string) { + t.Helper() + body, _ := json.Marshal(asked) + msg, err := conn.Request(CLISubject(node), body, 5*time.Second) + if err != nil { + t.Fatal(err) + } + var env struct { + Result json.RawMessage `json:"result"` + Error string `json:"error"` + } + if err := json.Unmarshal(msg.Data, &env); err != nil { + t.Fatalf("not an envelope: %s", msg.Data) + } + var a CLIAnswer + var raw map[string]any + _ = json.Unmarshal(env.Result, &a) + _ = json.Unmarshal(env.Result, &raw) + return a, raw, env.Error +} + +// **A line that outlasts the bus's window for an answer is followed to its answer, never lost** (review of ADR +// 0272): the first answer says it is running and names its call, and following the call by the same account on +// the same node gives what the line printed once it ends. Another account, or another node, is told no such call. +func TestALongMeshCLILineIsFollowedToItsAnswer(t *testing.T) { + was := AnswerWithin + AnswerWithin = 50 * time.Millisecond + t.Cleanup(func() { AnswerWithin = was }) + release := make(chan struct{}) + conn := cliBus(t, NewCallLog(), 4, func(ctx context.Context, node string, asked CLIAsked) CLIAnswer { + <-release + return CLIAnswer{Stdout: []byte("done on " + node), Terminal: true} + }) + _, first, failed := askCLI(t, conn, "laptop", CLIAsked{Line: []string{"push", "laptop"}, Account: "op", UID: 1000}) + call, _ := first["call"].(string) + if failed != "" || first["running"] != true || call == "" { + t.Fatalf("a long line was not answered as running with its call: %v %q", first, failed) + } + _, still, _ := askCLI(t, conn, "laptop", CLIAsked{Follow: call, Account: "op", UID: 1000}) + if still["running"] != true { + t.Fatalf("following a running line answered %v", still) + } + close(release) + deadline := time.Now().Add(5 * time.Second) + for { + a, raw, failed := askCLI(t, conn, "laptop", CLIAsked{Follow: call, Account: "op", UID: 1000}) + if failed != "" { + t.Fatalf("following answered %q", failed) + } + if raw["running"] != true { + if string(a.Stdout) != "done on laptop" || !a.Terminal { + t.Fatalf("followed to %+v", a) + } + break + } + if time.Now().After(deadline) { + t.Fatal("the line never finished for its follower") + } + time.Sleep(20 * time.Millisecond) + } + if _, _, failed := askCLI(t, conn, "laptop", CLIAsked{Follow: call, Account: "agent", UID: 1001}); failed == "" { + t.Fatal("another account followed the operator's line") + } + if _, _, failed := askCLI(t, conn, "desktop", CLIAsked{Follow: call, Account: "op", UID: 1000}); failed == "" { + t.Fatal("another node followed the line") + } +} + +// Lines running at once are bounded: one over the bound is answered busy at once, and nothing ran. +func TestMeshCLILinesAtOnceAreBounded(t *testing.T) { + was := AnswerWithin + AnswerWithin = 50 * time.Millisecond + t.Cleanup(func() { AnswerWithin = was }) + release := make(chan struct{}) + t.Cleanup(func() { close(release) }) + ran := make(chan struct{}, 4) + conn := cliBus(t, NewCallLog(), 1, func(context.Context, string, CLIAsked) CLIAnswer { + ran <- struct{}{} + <-release + return CLIAnswer{} + }) + if _, first, _ := askCLI(t, conn, "laptop", CLIAsked{Line: []string{"status"}, Account: "op"}); first["running"] != true { + t.Fatalf("the first line answered %v", first) + } + a, _, _ := askCLI(t, conn, "laptop", CLIAsked{Line: []string{"status"}, Account: "op"}) + if !strings.Contains(a.Refused, "busy") || a.Exit == 0 { + t.Fatalf("a line over the bound answered %+v", a) + } + if len(ran) != 1 { + t.Fatalf("%d lines ran, one was bound", len(ran)) + } +} + +// A mesh-cli line's record keeps its first word, never the rest of its line, and its answer is never kept on the +// bus, where `calls` answers anyone who may call the seat. +func TestAMeshCLIRecordKeepsNeitherItsLineNorItsAnswerOnTheBus(t *testing.T) { + l, a := NewCallLog(), newAnswers(t) + writes := make(chan Call, 4) + l.writes = writes + asked, _ := json.Marshal(CLIAsked{Line: []string{"settings", "set", "x", `{"password":"s3cret"}`}, Account: "op"}) + l.serveCall(CLISeat, "laptop", asked, "_INBOX.node.laptop.abcdefghijklmnopqrstuv", + func(context.Context, json.RawMessage) (any, error) { + return CLIAnswer{Stdout: []byte("token s3cret-join")}, nil + }, + a.respond, nil) + _ = a.only() + close(writes) + for c := range writes { + if strings.Contains(string(c.Args), "s3cret") || strings.Contains(string(c.Answer), "s3cret") { + t.Fatalf("the bus was sent %s / %s", c.Args, c.Answer) + } + if !strings.Contains(string(c.Args), "settings") || !strings.Contains(string(c.Args), `"op"`) { + t.Fatalf("the record does not say what was asked and by whom: %s", c.Args) + } + } +} + +// `calls` answers anyone who may call the seat, agents among them: a mesh-cli line's answer is never shown there, +// even from the memory of the controller that ran it; every other call's is (review of ADR 0272). +func TestCallsNeverShowsAMeshCLILinesAnswer(t *testing.T) { + l, a := NewCallLog(), newAnswers(t) + asked, _ := json.Marshal(CLIAsked{Line: []string{"token", "issue", "x"}, Account: "operator"}) + l.serveCall(CLISeat, "control", asked, "_INBOX.node.control.abcdefghijklmnopqrstuv", + func(context.Context, json.RawMessage) (any, error) { + return CLIAnswer{Stdout: []byte("s3cret-join")}, nil + }, + a.respond, nil) + _ = a.only() + recent, _ := l.Recent() + shown, found, err := l.Shown(recent[0].ID) + if err != nil || !found { + t.Fatalf("the line is not shown at all: %v %v", found, err) + } + if strings.Contains(string(shown.Answer), "s3cret-join") { + t.Fatalf("calls shows a mesh-cli line's answer: %s", shown.Answer) + } + b := newAnswers(t) + l.serveCall("mesh-controller", "status", nil, "_INBOX.x.abcdefghijklmnopqrstuv", + func(context.Context, json.RawMessage) (any, error) { return "all well", nil }, b.respond, nil) + _ = b.only() + recent, _ = l.Recent() + if shown, _, _ := l.Shown(recent[0].ID); !strings.Contains(string(shown.Answer), "all well") { + t.Fatalf("another call's answer is withheld: %s", shown.Answer) + } +}