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 } // 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(): 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 } } }