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).
This commit is contained in:
@@ -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 <module> <tool> [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
|
||||||
|
}
|
||||||
@@ -68,6 +68,8 @@ func run() error {
|
|||||||
return licenceCommand(ctx, args[1:])
|
return licenceCommand(ctx, args[1:])
|
||||||
case "rotate":
|
case "rotate":
|
||||||
return rotateCommand(ctx, args[1:])
|
return rotateCommand(ctx, args[1:])
|
||||||
|
case "ask":
|
||||||
|
return askCommand(ctx, args[1:])
|
||||||
case "builds":
|
case "builds":
|
||||||
return buildsCommand(ctx, args[1:])
|
return buildsCommand(ctx, args[1:])
|
||||||
case "pin":
|
case "pin":
|
||||||
@@ -171,6 +173,7 @@ func usage() {
|
|||||||
licence manager <name> <node> the node that holds a refreshable licence's refresh token
|
licence manager <name> <node> the node that holds a refreshable licence's refresh token
|
||||||
licence refresh <name> mint a new access token and seal it to every holder
|
licence refresh <name> mint a new access token and seal it to every holder
|
||||||
rotate <provision> [--consumer <n>] a new credential for every holder, both ends at once
|
rotate <provision> [--consumer <n>] a new credential for every holder, both ends at once
|
||||||
|
ask <module> <tool> [json] call one of a module's tools over the broker, and print its answer
|
||||||
pin <node> <provision> <from> which node this one gets a provision from
|
pin <node> <provision> <from> which node this one gets a provision from
|
||||||
unpin <node> <provision> put that question back
|
unpin <node> <provision> put that question back
|
||||||
plan <node> [--files|--json] what that node would run, and why
|
plan <node> [--files|--json] what that node would run, and why
|
||||||
|
|||||||
@@ -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 `<module>.<tool>`, 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
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user