Ask a seat once more when the controller refused as handing over (hq issue 289)
mesh/merge-gate pass: builds mesh-tools, node-tools → ace, g14, novox, shanks; no bus step; every machine composes with the change as it did without (4 of …
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery delivered
mesh/delivery-group group fix/a-verb-survives-the-handover delivered: every member is delivered

A controller being replaced refuses a call it cannot serve, marked retry:
handing-over; the one after it answers the same call.
This commit is contained in:
jochen
2026-10-07 02:13:17 +02:00
parent 01c98db0a4
commit 2f045e042f
2 changed files with 106 additions and 6 deletions
+33 -6
View File
@@ -578,28 +578,55 @@ func (c *Conn) Ask(key string, body any, on string) (Answered, error) {
if err != nil {
return Answered{}, err
}
answered, handingOver, err := c.askOnce(subject, data)
if !handingOver {
return answered, err
}
// **Refused because the controller is handing over: asked once more** (novox/hq issue 289). The
// controller being replaced did nothing and said so; the one after it answers the same call. Once:
// a second refusal is a controller that is not coming back, and the caller is told.
time.Sleep(HandoverPause)
answered, again, err := c.askOnce(subject, data)
if err != nil && again {
return Answered{}, fmt.Errorf("asked twice, %s apart, and both times refused as a handover: %w", HandoverPause, err)
}
return answered, err
}
// HandoverPause is how long a call refused as a handover waits before it is asked once more: the
// supervisor's restart of the controller, which stops the one being replaced and starts the next. A
// variable so a test need not wait.
var HandoverPause = 3 * time.Second
// RetryHandingOver is the mark the controller puts on a call it refused because it is being replaced
// (mesh-controller internal/link, `"retry": "handing-over"`).
const RetryHandingOver = "handing-over"
// askOnce asks once, and says whether the answer was a refusal marked to be asked again.
func (c *Conn) askOnce(subject string, data []byte) (Answered, bool, error) {
msg, err := c.nc.Request(subject, data, RequestTimeout)
if err != nil {
if errors.Is(err, nats.ErrNoResponders) {
return Answered{}, errors.New("503 no responders")
return Answered{}, false, errors.New("503 no responders")
}
if errors.Is(err, nats.ErrTimeout) {
return Answered{}, errors.New("timeout")
return Answered{}, false, errors.New("timeout")
}
return Answered{}, err
return Answered{}, false, err
}
var r struct {
Result json.RawMessage `json:"result"`
Error string `json:"error"`
Node string `json:"node"`
Retry string `json:"retry"`
}
if err := json.Unmarshal(msg.Data, &r); err != nil {
return Answered{}, err
return Answered{}, false, err
}
if r.Error != "" {
return Answered{}, errors.New(r.Error)
return Answered{}, r.Retry == RetryHandingOver, errors.New(r.Error)
}
return Answered{Result: r.Result, Node: r.Node}, nil
return Answered{Result: r.Result, Node: r.Node}, false, nil
}
// PublishAs emits an event as a module: published into JetStream and awaited, de-duplicated by its
+73
View File
@@ -0,0 +1,73 @@
package bus_test
import (
"encoding/json"
"strings"
"sync/atomic"
"testing"
"time"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-tools/node-tools/internal/bus"
"github.com/novox/mesh-tools/node-tools/internal/meshtest"
)
// **A call the controller refused because it is handing over is asked once more** (novox/hq issue 289):
// the controller being replaced did nothing and marked its refusal, and the one after it answers. A
// refusal not so marked is final, and a second handover refusal is said, not asked a third time.
func TestAHandoverRefusalIsAskedOnceMore(t *testing.T) {
was := bus.HandoverPause
bus.HandoverPause = 10 * time.Millisecond
t.Cleanup(func() { bus.HandoverPause = was })
url := meshtest.URL(t)
seat, err := nats.Connect(url)
if err != nil {
t.Fatal(err)
}
defer seat.Close()
c, err := bus.Connect(bus.Credential{URL: url, Module: "console", Node: "anchor"})
if err != nil {
t.Fatal(err)
}
defer c.Close()
serve := func(t *testing.T, verb string, answers ...map[string]any) *atomic.Int32 {
t.Helper()
var asked atomic.Int32
sub, err := seat.Subscribe(bus.SeatToolSubject("mesh-controller", verb, "", ""), func(m *nats.Msg) {
n := int(asked.Add(1)) - 1
if n >= len(answers) {
n = len(answers) - 1
}
body, _ := json.Marshal(answers[n])
_ = m.Respond(body)
})
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = sub.Unsubscribe() })
_ = seat.Flush()
return &asked
}
handingOver := map[string]any{"error": "the controller is handing over to the next one; nothing was done, ask again",
"retry": bus.RetryHandingOver}
asked := serve(t, "rotate", handingOver, map[string]any{"result": "rotated"})
got, err := c.Ask("seat:mesh-controller.rotate", map[string]any{}, "")
if err != nil || string(got.Result) != `"rotated"` || asked.Load() != 2 {
t.Fatalf("after a handover refusal: %s %v, asked %d times", got.Result, err, asked.Load())
}
asked = serve(t, "push", map[string]any{"error": "refused for its own reason"})
if _, err := c.Ask("seat:mesh-controller.push", map[string]any{}, ""); err == nil || asked.Load() != 1 {
t.Fatalf("an unmarked refusal was asked again: %v, asked %d times", err, asked.Load())
}
asked = serve(t, "status", handingOver)
_, err = c.Ask("seat:mesh-controller.status", map[string]any{}, "")
if err == nil || asked.Load() != 2 || !strings.Contains(err.Error(), "both times refused as a handover") {
t.Fatalf("two handover refusals: %v, asked %d times", err, asked.Load())
}
}