From 70705ffe457c4df5e34c1e6835905a06361f7217 Mon Sep 17 00:00:00 2001 From: jochen Date: Wed, 30 Sep 2026 17:59:23 +0200 Subject: [PATCH] A holder binds its seat's tools when it may, not only when it starts 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. --- internal/link/seattools.go | 95 +++++++++++++++++++++++++++++--------- 1 file changed, 73 insertions(+), 22 deletions(-) diff --git a/internal/link/seattools.go b/internal/link/seattools.go index 11c7f75..f579c53 100644 --- a/internal/link/seattools.go +++ b/internal/link/seattools.go @@ -29,11 +29,23 @@ func SeatToolSubject(seat, verb string) string { return "mesh.seat." + seat + ". // 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() } @@ -41,36 +53,75 @@ func (b OverNATS) ServeSeatTools(seat string, handlers map[string]ToolHandler, l 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) - } - }() - }) + 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 + } +} -- 2.54.0