Files
mesh-controller/internal/link/ask.go
T
jschoubben 973cda5d76 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).
2026-09-21 20:34:03 +02:00

87 lines
3.1 KiB
Go

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
}
}
}