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