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 // 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. func (b OverNATS) ServeSeatTools(seat string, handlers map[string]ToolHandler, logger *log.Logger) (func(), error) { var subs []*nats.Subscription stop := func() { for _, s := range subs { _ = s.Unsubscribe() } } for verb, handle := range handlers { verb, handle := verb, handle subject := SeatToolSubject(seat, verb) sub, err := 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) } }() }) if err != nil { stop() return nil, fmt.Errorf("serving %s: %w", subject, err) } subs = append(subs, sub) } if logger != nil { logger.Printf("serving %d tool(s) of the %s seat", len(handlers), seat) } return stop, nil }