mesh-cli lines are calls, followed to their answer; the terminal does not follow assign (review of hq ADR 0272)
- A line ran past the bus's one-minute window for an answer and its answer was refused; mesh-cli then said nothing ran of a line that had. Every line is now a call (calls.go): answered within AnswerWithin that it is running, with its call, and followed by the same account on the same node until it ends. Its record keeps the command word only, and its answer stays in the memory of the controller that ran it — never on the bus, never in `calls`, which answers anyone who may call the seat. Each line is said in the journal with its call, who asked where, and how. - Lines running at once are bounded (8); one over is answered busy. - An ordinary `settings` line is composed as the settings verb composes its own and meets that verb's refusals, the terminal-only keys among them, instead of the generic command's blanket refusal. - The controller's module is assigned and unassigned at the terminal only: where it runs is the control-node, whose operator is the terminal. - mesh.control.*.cli is in the writers table as the machine's own.
This commit is contained in:
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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 <module> <values> [--replace] [--node <node>], settings clear <module> " +
|
||||
"[--node <node>], settings show <module> [--history] [--node <node>], or settings preferences [<module>] " +
|
||||
"[--node <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)
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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.<its own>.>`).
|
||||
Subjects: []string{"mesh.control.*.report", "mesh.control.*.health"}, Writes: ownMachine},
|
||||
// machine, inside the grant it already had (`mesh.control.<its own>.>`). 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",
|
||||
|
||||
@@ -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.>"}},
|
||||
|
||||
@@ -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) }) })
|
||||
|
||||
+108
-22
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user