One bus: the AMQP transport is gone from the controller
The mesh runs on the seat's bus alone (novox/hq ADR 0131, design 28 task 5.5). The old transport's consume loop, build request, tool ask, management API and account scoping are deleted, and the bus switch with them; the controller connects to the broker seat and to nothing else. The store-window tests keep their assertions on a bus-less fake, and the tests that only made sense for the old transport's in-memory holding go with it.
This commit is contained in:
+22
-60
@@ -3,16 +3,13 @@ package link
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
"github.com/nats-io/nats.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"`
|
||||
@@ -27,70 +24,35 @@ type Answer struct {
|
||||
// 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,
|
||||
// The answer comes back on the asker's own inbox, which only the asker may read; the serving
|
||||
// module answers there and nowhere else. The bus refuses a request nothing serves at once, so a
|
||||
// module that is down or a tool that does not exist is said now rather than after the whole wait.
|
||||
func Ask(ctx context.Context, bus Bus, 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
|
||||
}
|
||||
|
||||
// Mandatory, so a request nothing consumes comes straight back: a module that is down, or a
|
||||
// tool that does not exist, is said at once rather than after the whole wait.
|
||||
returned := channel.NotifyReturn(make(chan amqp.Return, 1))
|
||||
id := fmt.Sprintf("ask-%d", time.Now().UnixNano())
|
||||
key := module + "." + tool
|
||||
if err := channel.PublishWithContext(ctx, RPCExchange, key, true, 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 back := <-returned:
|
||||
if back.CorrelationId == id {
|
||||
return Answer{}, fmt.Errorf(
|
||||
"nothing serves %s: no runtime has bound %q on the broker. The module is not "+
|
||||
"assigned, its runtime is not up, or it serves no such tool — `status` "+
|
||||
"says whether the machine carrying it has applied", module, key)
|
||||
}
|
||||
case <-waiting.Done():
|
||||
reply, err := bus.AskTool(ctx, module, tool, args, timeout)
|
||||
if err != nil {
|
||||
if errors.Is(err, nats.ErrNoResponders) {
|
||||
return Answer{}, fmt.Errorf(
|
||||
"nothing serves %s: no runtime has bound %q on the bus. The module is not "+
|
||||
"assigned, its runtime is not up, or it serves no such tool — `status` says "+
|
||||
"whether the machine carrying it has applied", module, key)
|
||||
}
|
||||
if errors.Is(err, context.DeadlineExceeded) || errors.Is(err, nats.ErrTimeout) {
|
||||
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",
|
||||
"the bus — `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
|
||||
}
|
||||
return Answer{}, fmt.Errorf("cannot ask %s: %w", key, err)
|
||||
}
|
||||
var answer Answer
|
||||
if err := json.Unmarshal(reply, &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