Merge pull request 'Answer mesh-cli: the control-node's operator is the terminal, everyone else is not (hq ADR 0272)' (#176) from feat/272-answer-mesh-cli into main
This commit was merged in pull request #176.
This commit is contained in:
@@ -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) }) })
|
||||
|
||||
@@ -0,0 +1,216 @@
|
||||
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"`
|
||||
}
|
||||
|
||||
// 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
|
||||
}
|
||||
@@ -0,0 +1,265 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/base64"
|
||||
"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
|
||||
// same list (mesh-host internal/meshcli, TestTheRequestToTheControllerKeepsItsFieldNames).
|
||||
func TestTheMeshCLIRequestAndAnswerKeepTheirFieldNames(t *testing.T) {
|
||||
body, _ := json.Marshal(CLIAsked{Line: []string{"status"}, Account: "a", UID: 1})
|
||||
if got := keysIn(t, body); got != "account line uid" {
|
||||
t.Fatalf("the request's fields are %q", got)
|
||||
}
|
||||
body, _ = json.Marshal(CLIAnswer{Stdout: []byte("o"), Stderr: []byte("e"), Exit: 1, Terminal: true, Why: "w",
|
||||
Refused: "r", Cut: true})
|
||||
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 {
|
||||
t.Helper()
|
||||
var m map[string]any
|
||||
if err := json.Unmarshal(body, &m); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var keys []string
|
||||
for k := range m {
|
||||
keys = append(keys, k)
|
||||
}
|
||||
sort.Strings(keys)
|
||||
return strings.Join(keys, " ")
|
||||
}
|
||||
|
||||
// A node's mesh-cli subject names the node, and no other subject does.
|
||||
func TestTheMeshCLISubjectNamesItsNode(t *testing.T) {
|
||||
if node, ok := CLINode(CLISubject("laptop")); !ok || node != "laptop" {
|
||||
t.Fatalf("CLINode(%q) = %q %v", CLISubject("laptop"), node, ok)
|
||||
}
|
||||
for _, s := range []string{"mesh.control.laptop.report", "mesh.control..cli", "mesh.control.a.b.cli", "mesh.node.a.cli"} {
|
||||
if _, ok := CLINode(s); ok {
|
||||
t.Fatalf("%s was read as a mesh-cli subject", s)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// An answer larger than one bus message is cut to fit, and the cut is said.
|
||||
func TestALargeMeshCLIAnswerIsCutAndSaysSo(t *testing.T) {
|
||||
a := CLIAnswer{Stdout: []byte(strings.Repeat("o", 200000)), Stderr: []byte(strings.Repeat("e", 1000)), Why: "the terminal"}
|
||||
body := FitCLIAnswer(a, 64<<10)
|
||||
if len(body) > 64<<10 {
|
||||
t.Fatalf("the answer is %d bytes, over the bound", len(body))
|
||||
}
|
||||
var got CLIAnswer
|
||||
if err := json.Unmarshal(body, &got); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !got.Cut || got.Why != "the terminal" || len(got.Stdout) == 0 || string(got.Stderr) != strings.Repeat("e", 1000) {
|
||||
t.Fatalf("the cut answer is %+v", got)
|
||||
}
|
||||
if small := FitCLIAnswer(CLIAnswer{Stdout: []byte("ok")}, 64<<10); strings.Contains(string(small), `"cut"`) {
|
||||
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("s3cret-join")}, nil
|
||||
},
|
||||
a.respond, nil)
|
||||
_ = a.only()
|
||||
close(writes)
|
||||
for c := range writes {
|
||||
if carries(c.Args, "s3cret") || carries(c.Answer, "s3cret-join") {
|
||||
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 carries(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)
|
||||
}
|
||||
}
|
||||
|
||||
// carries says whether a record holds a secret as text or as the base64 JSON writes bytes in: a line's output is
|
||||
// bytes, so a search for its text alone finds nothing whatever the record keeps (review of ADR 0272).
|
||||
func carries(body []byte, secret string) bool {
|
||||
return strings.Contains(string(body), secret) ||
|
||||
strings.Contains(string(body), base64.StdEncoding.EncodeToString([]byte(secret)))
|
||||
}
|
||||
|
||||
// The two tests above hold what they claim: each fails when the protection it names is taken away.
|
||||
func TestTheWithholdingTestsHoldSomething(t *testing.T) {
|
||||
// The answer as a mesh-cli line's record would carry it, were it kept: the search finds it.
|
||||
body, _ := json.Marshal(map[string]any{"result": CLIAnswer{Stdout: []byte("s3cret-join")}})
|
||||
if !carries(body, "s3cret-join") {
|
||||
t.Fatalf("a record carrying the answer is not found carrying it: %s", body)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user