From 2f045e042f2513075647b9821110c8b217fdcfdb Mon Sep 17 00:00:00 2001 From: jochen Date: Wed, 7 Oct 2026 02:13:17 +0200 Subject: [PATCH] Ask a seat once more when the controller refused as handing over (hq issue 289) A controller being replaced refuses a call it cannot serve, marked retry: handing-over; the one after it answers the same call. --- node-tools/internal/bus/bus.go | 39 +++++++++++-- node-tools/internal/bus/handover_test.go | 73 ++++++++++++++++++++++++ 2 files changed, 106 insertions(+), 6 deletions(-) create mode 100644 node-tools/internal/bus/handover_test.go diff --git a/node-tools/internal/bus/bus.go b/node-tools/internal/bus/bus.go index 64f543d..0013c71 100644 --- a/node-tools/internal/bus/bus.go +++ b/node-tools/internal/bus/bus.go @@ -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 diff --git a/node-tools/internal/bus/handover_test.go b/node-tools/internal/bus/handover_test.go new file mode 100644 index 0000000..a3478f4 --- /dev/null +++ b/node-tools/internal/bus/handover_test.go @@ -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()) + } +}