From 973cda5d769318c588f023fd0d95619d9353139c Mon Sep 17 00:00:00 2001 From: jochen Date: Mon, 21 Sep 2026 20:34:03 +0200 Subject: [PATCH] ask: the control plane calls a module's tool and prints its answer A module serves tools under an account scoped to exactly that, and nothing else in the mesh held an account that could ask one. The control plane does: ask publishes on the RPC exchange with a private reply queue bound under its own name, checks the correlation, prints the answer, and exits non-zero for a tool that answered with an error or a module that never answered (novox/hq 04-ISSUES/049, ADR 0095). --- cmd/mesh-controller/ask.go | 62 ++++++++++++++++++++++++++ cmd/mesh-controller/main.go | 3 ++ internal/link/ask.go | 86 +++++++++++++++++++++++++++++++++++++ 3 files changed, 151 insertions(+) create mode 100644 cmd/mesh-controller/ask.go create mode 100644 internal/link/ask.go diff --git a/cmd/mesh-controller/ask.go b/cmd/mesh-controller/ask.go new file mode 100644 index 0000000..001d5e2 --- /dev/null +++ b/cmd/mesh-controller/ask.go @@ -0,0 +1,62 @@ +package main + +import ( + "context" + "encoding/json" + "errors" + "flag" + "fmt" + "os" + "time" + + "github.com/novox/mesh-controller/internal/link" +) + +// ask calls one of a module's tools, through the control plane's own broker connection. +// +// A module serves tools under an account scoped to exactly that (novox/hq ADR 0047), and nothing +// else in the mesh held an account that could ask one — not an operator at a terminal, not an agent +// acting for one (novox/hq 04-ISSUES/049). The control plane does, so it is the way in: one +// process, one connection, one place a question can be seen to have been asked (ADR 0095). +func askCommand(ctx context.Context, args []string) error { + positionals, flags := split(args) + set := flag.NewFlagSet("ask", flag.ContinueOnError) + wait := set.Duration("wait", 60*time.Second, "how long to wait for the module's answer") + if err := set.Parse(flags); err != nil { + return err + } + if len(positionals) < 2 || len(positionals) > 3 { + return errors.New("ask [json arguments] [--wait 60s]") + } + module, tool := positionals[0], positionals[1] + var arguments json.RawMessage + if len(positionals) == 3 { + if !json.Valid([]byte(positionals[2])) { + return fmt.Errorf("the arguments are not JSON: %s", positionals[2]) + } + arguments = json.RawMessage(positionals[2]) + } + + server, err := link.Connect(nil, nil) + if err != nil { + return err + } + defer server.Close() + + answer, err := link.Ask(ctx, server.Channel(), module, tool, arguments, *wait) + if err != nil { + return err + } + // The answer as the module gave it, to standard output, for a person or a program. A tool + // that answered with an error has still answered: printed the same way, and the exit status + // says which. + body, err := json.Marshal(answer) + if err != nil { + return err + } + fmt.Fprintln(os.Stdout, string(body)) + if answer.Error != "" { + return fmt.Errorf("%s.%s answered with an error: %s", module, tool, answer.Error) + } + return nil +} diff --git a/cmd/mesh-controller/main.go b/cmd/mesh-controller/main.go index 08bf58e..70b6bfb 100644 --- a/cmd/mesh-controller/main.go +++ b/cmd/mesh-controller/main.go @@ -68,6 +68,8 @@ func run() error { return licenceCommand(ctx, args[1:]) case "rotate": return rotateCommand(ctx, args[1:]) + case "ask": + return askCommand(ctx, args[1:]) case "builds": return buildsCommand(ctx, args[1:]) case "pin": @@ -171,6 +173,7 @@ func usage() { licence manager the node that holds a refreshable licence's refresh token licence refresh mint a new access token and seal it to every holder rotate [--consumer ] a new credential for every holder, both ends at once + ask [json] call one of a module's tools over the broker, and print its answer pin which node this one gets a provision from unpin put that question back plan [--files|--json] what that node would run, and why diff --git a/internal/link/ask.go b/internal/link/ask.go new file mode 100644 index 0000000..df29a8b --- /dev/null +++ b/internal/link/ask.go @@ -0,0 +1,86 @@ +package link + +import ( + "context" + "encoding/json" + "fmt" + "time" + + amqp "github.com/rabbitmq/amqp091-go" +) + +// RPCExchange is where a module's tools are asked over the broker, keyed `.`, and +// where the answer comes back, keyed by the asker's reply queue (novox/hq ADR 0047). +const RPCExchange = "mesh.rpc" + +// Answer is what a module's tool replies: one of the two, never both. +type Answer struct { + Result json.RawMessage `json:"result,omitempty"` + Error string `json:"error,omitempty"` +} + +// Ask calls one of a module's tools over the broker and waits for its answer. +// +// **The control plane is the way in** (novox/hq 04-ISSUES/049, ADR 0095). A module's broker +// account is scoped to what it emits, consumes and serves, and a tool call needs a reply queue the +// caller creates and a publish to the serving module's request key — which no module's scope +// grants, and should not. The control plane already holds a connection that may, so a person or +// an agent asks through it, and every question passes one process where an audit belongs. +// +// The reply queue is the caller's own, server-named and exclusive, bound to the RPC exchange under +// its own name: a serving module answers through that exchange and never the default one, whose +// permission is per exchange rather than per queue. The correlation is checked rather than +// assumed, as every RPC here is. +func Ask(ctx context.Context, channel *amqp.Channel, module, tool string, args json.RawMessage, + timeout time.Duration) (Answer, error) { + + if len(args) == 0 { + args = json.RawMessage(`{}`) + } + replies, err := channel.QueueDeclare("", false, true, true, false, nil) + if err != nil { + return Answer{}, err + } + if err := channel.QueueBind(replies.Name, replies.Name, RPCExchange, false, nil); err != nil { + return Answer{}, fmt.Errorf("cannot bind a reply queue to %s: %w", RPCExchange, err) + } + answers, err := channel.ConsumeWithContext(ctx, replies.Name, "", true, true, false, false, nil) + if err != nil { + return Answer{}, err + } + + id := fmt.Sprintf("ask-%d", time.Now().UnixNano()) + key := module + "." + tool + if err := channel.PublishWithContext(ctx, RPCExchange, key, false, false, amqp.Publishing{ + ContentType: "application/json", + CorrelationId: id, + ReplyTo: replies.Name, + Body: args, + }); err != nil { + return Answer{}, fmt.Errorf("cannot ask %s: %w", key, err) + } + + waiting, cancel := context.WithTimeout(ctx, timeout) + defer cancel() + for { + select { + case <-waiting.Done(): + return Answer{}, fmt.Errorf( + "%s did not answer within %s. Its runtime serves %q when it is up and has bound "+ + "the broker — `status` says whether the machine carrying it has applied", + module, timeout, key) + case delivery, ok := <-answers: + if !ok { + return Answer{}, fmt.Errorf("the connection closed while waiting for %s", key) + } + if delivery.CorrelationId != id { + continue + } + var answer Answer + if err := json.Unmarshal(delivery.Body, &answer); err != nil { + return Answer{}, fmt.Errorf("%s answered with something unreadable: %w", key, err) + } + return answer, nil + } + } +}