A push outlasted the console's 30s wait and, when it sent the bus its changed user list, the broker's reload forgot the reply it may send: the push happened and its caller was told it did not answer. Calls now answer in full or as running with an id, a push answers before it sends, refused answers are recorded on their call, and 'calls' reads them back.
126 lines
5.2 KiB
Go
126 lines
5.2 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 call. Its caller is answered within AnswerWithin either way; this is how
|
|
// long the call itself may run before it is stopped.
|
|
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) {
|
|
return b.serveTools(seat, func(verb string) string { return SeatToolSubject(seat, verb) }, handlers, logger)
|
|
}
|
|
|
|
// ServeNodeSeatTools is ServeSeatTools for one machine's holder of a node-scoped seat (novox/hq ADR
|
|
// 0159, ADR 0219): each verb on the seat's subject for this machine and no other, so a call names
|
|
// the machine it is for and only that machine's holder answers it.
|
|
func (b OverNATS) ServeNodeSeatTools(seat, node string, handlers map[string]ToolHandler, logger *log.Logger) (func(), error) {
|
|
return b.serveTools(seat, func(verb string) string { return NodeSeatToolSubject(seat, verb, node) }, handlers, logger)
|
|
}
|
|
|
|
func (b OverNATS) serveTools(seat string, subjectOf func(string) 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()
|
|
}
|
|
}
|
|
// A refused answer is recorded against its call, not only printed by the library.
|
|
Calls.WatchRefusals(b.Conn, logger)
|
|
for verb, handle := range handlers {
|
|
verb, handle := verb, handle
|
|
subject := subjectOf(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. Answered once, within
|
|
// AnswerWithin, and kept (novox/hq issue 265).
|
|
go Calls.serveCall(seat, verb, json.RawMessage(msg.Data), msg.Reply, handle, msg.Respond, logger)
|
|
})
|
|
}
|
|
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
|
|
}
|
|
}
|