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.
This commit is contained in:
+73
-22
@@ -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.
|
// answers in seconds; anything that has not in this long is said to have not answered.
|
||||||
const HandlerTimeout = 5 * time.Minute
|
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
|
// 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 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) {
|
func (b OverNATS) ServeSeatTools(seat string, handlers map[string]ToolHandler, logger *log.Logger) (func(), error) {
|
||||||
var subs []*nats.Subscription
|
var subs []*nats.Subscription
|
||||||
|
done := make(chan struct{})
|
||||||
stop := func() {
|
stop := func() {
|
||||||
|
close(done)
|
||||||
for _, s := range subs {
|
for _, s := range subs {
|
||||||
_ = s.Unsubscribe()
|
_ = s.Unsubscribe()
|
||||||
}
|
}
|
||||||
@@ -41,36 +53,75 @@ func (b OverNATS) ServeSeatTools(seat string, handlers map[string]ToolHandler, l
|
|||||||
for verb, handle := range handlers {
|
for verb, handle := range handlers {
|
||||||
verb, handle := verb, handle
|
verb, handle := verb, handle
|
||||||
subject := SeatToolSubject(seat, verb)
|
subject := SeatToolSubject(seat, verb)
|
||||||
sub, err := b.Conn.QueueSubscribe(subject, "seat."+seat, func(msg *nats.Msg) {
|
bind := func() (*nats.Subscription, error) {
|
||||||
// Its own goroutine per call: a slow `push` must not hold up a `status` asked beside it,
|
return b.Conn.QueueSubscribe(subject, "seat."+seat, func(msg *nats.Msg) {
|
||||||
// and the library would otherwise run handlers one after another.
|
// Its own goroutine per call: a slow `push` must not hold up a `status` asked beside it,
|
||||||
go func() {
|
// and the library would otherwise run handlers one after another.
|
||||||
ctx, cancel := context.WithTimeout(context.Background(), HandlerTimeout)
|
go func() {
|
||||||
defer cancel()
|
ctx, cancel := context.WithTimeout(context.Background(), HandlerTimeout)
|
||||||
args := json.RawMessage(msg.Data)
|
defer cancel()
|
||||||
if len(args) == 0 {
|
args := json.RawMessage(msg.Data)
|
||||||
args = json.RawMessage(`{}`)
|
if len(args) == 0 {
|
||||||
}
|
args = json.RawMessage(`{}`)
|
||||||
var reply []byte
|
}
|
||||||
result, err := handle(ctx, args)
|
var reply []byte
|
||||||
if err != nil {
|
result, err := handle(ctx, args)
|
||||||
reply, _ = json.Marshal(map[string]any{"error": err.Error()})
|
if err != nil {
|
||||||
} else if reply, err = json.Marshal(map[string]any{"result": result}); err != nil {
|
reply, _ = json.Marshal(map[string]any{"error": err.Error()})
|
||||||
reply, _ = json.Marshal(map[string]any{"error": "the answer could not be written as JSON: " + 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 := msg.Respond(reply); err != nil && logger != nil {
|
||||||
}
|
logger.Printf("%s: could not answer: %v", subject, err)
|
||||||
}()
|
}
|
||||||
})
|
}()
|
||||||
|
})
|
||||||
|
}
|
||||||
|
sub, err := bind()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
stop()
|
stop()
|
||||||
return nil, fmt.Errorf("serving %s: %w", subject, err)
|
return nil, fmt.Errorf("serving %s: %w", subject, err)
|
||||||
}
|
}
|
||||||
subs = append(subs, sub)
|
subs = append(subs, sub)
|
||||||
|
go keepBound(sub, bind, subject, done, logger)
|
||||||
}
|
}
|
||||||
if logger != nil {
|
if logger != nil {
|
||||||
logger.Printf("serving %d tool(s) of the %s seat", len(handlers), seat)
|
logger.Printf("serving %d tool(s) of the %s seat", len(handlers), seat)
|
||||||
}
|
}
|
||||||
return stop, nil
|
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
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user