The grant is a line in the bus's user list the controller itself composes and a push delivers, so the first controller to serve its seat started before the list named it and every subscription was refused for good (2026-09-30). A refused subscription is retried until it holds.
128 lines
4.9 KiB
Go
128 lines
4.9 KiB
Go
package link
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"log"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
)
|
|
|
|
// A role's tools, served by its holder (novox/hq ADR 0132, ADR 0154).
|
|
//
|
|
// The mesh's own verbs — `status`, `push`, `assign` — are the mesh-controller seat's tools, and the
|
|
// control plane is that seat's holder. So it answers them here, on the seat's subjects, the way a
|
|
// module's runtime answers a module's: one request, one reply on the asker's own inbox, `{result}` or
|
|
// `{error}`. Nothing about the transport is the command's business; a handler is a function of its
|
|
// arguments and gets the same answer the command line prints.
|
|
|
|
// ToolHandler answers one call of a role's tool. What it returns is marshalled as the result; an
|
|
// error is the tool answering with one, which is an answer and not a timeout.
|
|
type ToolHandler func(ctx context.Context, args json.RawMessage) (any, error)
|
|
|
|
// SeatToolSubject is where a mesh-scoped seat's tool is asked (design 33 §4).
|
|
func SeatToolSubject(seat, verb string) string { return "mesh.seat." + seat + ".tool." + verb }
|
|
|
|
// HandlerTimeout bounds one answer. A verb that runs a command — a push, a build with no wait —
|
|
// answers in seconds; anything that has not in this long is said to have not answered.
|
|
const HandlerTimeout = 5 * time.Minute
|
|
|
|
// RebindAfter is how long a refused subscription waits before it is tried again.
|
|
const RebindAfter = 30 * time.Second
|
|
|
|
// ServeSeatTools binds every handler on its seat's subject until stopped. A queue group per seat, so
|
|
// a second holder during a handover shares the calls rather than both answering one.
|
|
//
|
|
// **A holder binds when it may, not only when it starts.** The grant that lets the controller
|
|
// subscribe its seat's tools is a line in the bus's user list, and that list is composed by the
|
|
// controller and delivered to the broker's machine by a push — so the first controller to serve
|
|
// these started before the list named them, the server refused every subscription, and nothing
|
|
// tried again (2026-09-30). A refused subscription is therefore retried until it holds: the server
|
|
// says so asynchronously and invalidates the subscription, which is what is checked.
|
|
func (b OverNATS) ServeSeatTools(seat string, handlers map[string]ToolHandler, logger *log.Logger) (func(), error) {
|
|
var subs []*nats.Subscription
|
|
done := make(chan struct{})
|
|
stop := func() {
|
|
close(done)
|
|
for _, s := range subs {
|
|
_ = s.Unsubscribe()
|
|
}
|
|
}
|
|
for verb, handle := range handlers {
|
|
verb, handle := verb, handle
|
|
subject := SeatToolSubject(seat, verb)
|
|
bind := func() (*nats.Subscription, error) {
|
|
return b.Conn.QueueSubscribe(subject, "seat."+seat, func(msg *nats.Msg) {
|
|
// Its own goroutine per call: a slow `push` must not hold up a `status` asked beside it,
|
|
// and the library would otherwise run handlers one after another.
|
|
go func() {
|
|
ctx, cancel := context.WithTimeout(context.Background(), HandlerTimeout)
|
|
defer cancel()
|
|
args := json.RawMessage(msg.Data)
|
|
if len(args) == 0 {
|
|
args = json.RawMessage(`{}`)
|
|
}
|
|
var reply []byte
|
|
result, err := handle(ctx, args)
|
|
if err != nil {
|
|
reply, _ = json.Marshal(map[string]any{"error": err.Error()})
|
|
} else if reply, err = json.Marshal(map[string]any{"result": result}); err != nil {
|
|
reply, _ = json.Marshal(map[string]any{"error": "the answer could not be written as JSON: " + err.Error()})
|
|
}
|
|
if err := msg.Respond(reply); err != nil && logger != nil {
|
|
logger.Printf("%s: could not answer: %v", subject, err)
|
|
}
|
|
}()
|
|
})
|
|
}
|
|
sub, err := bind()
|
|
if err != nil {
|
|
stop()
|
|
return nil, fmt.Errorf("serving %s: %w", subject, err)
|
|
}
|
|
subs = append(subs, sub)
|
|
go keepBound(sub, bind, subject, done, logger)
|
|
}
|
|
if logger != nil {
|
|
logger.Printf("serving %d tool(s) of the %s seat", len(handlers), seat)
|
|
}
|
|
return stop, nil
|
|
}
|
|
|
|
// keepBound watches one subscription and re-binds it after the server refused it, until stopped.
|
|
// A subscription the server refused is invalid a moment after it was made; one it accepted stays
|
|
// valid. Checked rather than hooked, because the connection's error handler belongs to whoever
|
|
// dialled and a second one would replace it.
|
|
func keepBound(sub *nats.Subscription, bind func() (*nats.Subscription, error), subject string,
|
|
done <-chan struct{}, logger *log.Logger) {
|
|
current := sub
|
|
for {
|
|
select {
|
|
case <-done:
|
|
return
|
|
case <-time.After(3 * time.Second):
|
|
}
|
|
if current.IsValid() {
|
|
// Settled; from here a lost subscription is a lost connection, which the client
|
|
// restores itself with every subscription it holds.
|
|
return
|
|
}
|
|
if logger != nil {
|
|
logger.Printf("%s: the bus refused the subscription; trying again in %s — the grant "+
|
|
"arrives with the next push to the machine holding mesh-broker", subject, RebindAfter)
|
|
}
|
|
select {
|
|
case <-done:
|
|
return
|
|
case <-time.After(RebindAfter):
|
|
}
|
|
again, err := bind()
|
|
if err != nil {
|
|
continue
|
|
}
|
|
current = again
|
|
}
|
|
}
|