From 85810ccc8cfd8e25e5677da197ce5c7fc6a42a09 Mon Sep 17 00:00:00 2001 From: jochen Date: Sun, 4 Oct 2026 01:20:44 +0200 Subject: [PATCH] Ask the bus again for a subscription it refused (novox/hq issue 222) A push sends the machine that needs a grant and the bus's machine in the same breath, and the runtime can subscribe the moment before the bus reloads its user list. Refused once, the subscription stayed dead until some later membership re-served it, and a newly assigned module ran unreachable. A refused subject the runtime answers is now asked for again for about five minutes. --- node-tools/internal/bus/bus.go | 116 +++++++++++++++++++++++- node-tools/internal/bus/refused_test.go | 86 ++++++++++++++++++ 2 files changed, 199 insertions(+), 3 deletions(-) create mode 100644 node-tools/internal/bus/refused_test.go diff --git a/node-tools/internal/bus/bus.go b/node-tools/internal/bus/bus.go index adb896b..d1a7af3 100644 --- a/node-tools/internal/bus/bus.go +++ b/node-tools/internal/bus/bus.go @@ -157,6 +157,97 @@ type Conn struct { onNew []func(Membership) subs []*nats.Subscription Logf func(format string, args ...any) + // answering is every subject this runtime answers on, so a subscription the bus refused can be + // asked for again (novox/hq issue 222). + answering map[string]*answered + // RetryAfter is how long after a refusal a subscription is asked for again, by attempt; past the + // last, it is given up and said. A variable so a test need not wait minutes. + RetryAfter []time.Duration +} + +// answered is one subject the runtime answers on: what to subscribe again with, and how often it was +// refused. +type answered struct { + subject, queue string + cb nats.MsgHandler + sub *nats.Subscription + refusals int + stopped bool +} + +// retryAfter is the default wait before each new attempt: a push sends the bus's user list in the +// same breath as the machine that needs it, and the bus reloads it a moment later — five minutes +// covers a slow push, and a grant that never comes is said and left. +var retryAfter = []time.Duration{2 * time.Second, 5 * time.Second, 10 * time.Second, 20 * time.Second, + 30 * time.Second, 60 * time.Second, 60 * time.Second, 120 * time.Second} + +var ( + refusedSubject = regexp.MustCompile(`Subscription to "(\S+)"`) + refusedQueue = regexp.MustCompile(`using queue "(\S+)"`) +) + +func answeringKey(subject, queue string) string { return subject + "\x00" + queue } + +// refused is the bus saying no to a subscription. **Asked again, not given up** (novox/hq issue 222): +// what an account may answer is the bus's user list, written on the bus's machine, and a push that +// assigns a module somewhere sends that machine and the bus's in the same breath — the runtime can +// subscribe in the moment before the bus has reloaded. Refused once, the subscription stayed dead +// until some later membership happened to re-serve it, and the module ran unreachable meanwhile. +func (c *Conn) refused(err error) { + if !errors.Is(err, nats.ErrPermissionViolation) { + return + } + m := refusedSubject.FindStringSubmatch(err.Error()) + if len(m) < 2 { + return + } + queue := "" + if q := refusedQueue.FindStringSubmatch(err.Error()); len(q) >= 2 { + queue = q[1] + } + c.mu.Lock() + a := c.answering[answeringKey(m[1], queue)] + if a == nil || a.stopped { + c.mu.Unlock() + return + } + waits := c.RetryAfter + if waits == nil { + waits = retryAfter + } + if a.refusals >= len(waits) { + c.mu.Unlock() + c.Logf("[mesh-tools] the bus still refuses %s after %d attempts; not served here until the mesh issues it again", a.subject, a.refusals) + return + } + wait := waits[a.refusals] + a.refusals++ + attempt := a.refusals + c.mu.Unlock() + c.Logf("[mesh-tools] the bus refused %s; asking again in %s (attempt %d)", a.subject, wait, attempt) + time.AfterFunc(wait, func() { + c.mu.Lock() + defer c.mu.Unlock() + if a.stopped { + return + } + if a.sub != nil { + _ = a.sub.Unsubscribe() + } + var sub *nats.Subscription + var err error + if a.queue != "" { + sub, err = c.nc.QueueSubscribe(a.subject, a.queue, a.cb) + } else { + sub, err = c.nc.Subscribe(a.subject, a.cb) + } + if err != nil { + c.Logf("[mesh-tools] asking again for %s failed: %v", a.subject, err) + return + } + a.sub = sub + c.subs = append(c.subs, sub) + }) } // Connect dials the bus as the credential's module. A module's subjects come from its credential, @@ -170,8 +261,14 @@ func Connect(cred Credential) (*Conn, error) { if node == "" { node = "?" } + var c *Conn opts := []nats.Option{ nats.Name(node + "." + cred.Module), + nats.ErrorHandler(func(_ *nats.Conn, _ *nats.Subscription, err error) { + if c != nil { + c.refused(err) + } + }), // Reconnect forever: the bus restarting is an upgrade, not a reason to exit. nats.MaxReconnects(-1), } @@ -192,8 +289,8 @@ func Connect(cred Credential) (*Conn, error) { nc.Close() return nil, err } - c := &Conn{nc: nc, js: js, self: cred.Module, node: cred.Node, cred: cred, - issued: map[string]*Membership{}, Logf: log.Printf} + c = &Conn{nc: nc, js: js, self: cred.Module, node: cred.Node, cred: cred, + issued: map[string]*Membership{}, Logf: log.Printf, answering: map[string]*answered{}} c.Follow(cred.Module) return c, nil } @@ -362,7 +459,20 @@ func (c *Conn) answerOn(subject, queue string, h Handler) (func(), error) { return func() {}, err } c.track(sub) - return func() { _ = sub.Unsubscribe() }, nil + a := &answered{subject: subject, queue: queue, cb: cb, sub: sub} + c.mu.Lock() + c.answering[answeringKey(subject, queue)] = a + c.mu.Unlock() + return func() { + c.mu.Lock() + a.stopped = true + if c.answering[answeringKey(subject, queue)] == a { + delete(c.answering, answeringKey(subject, queue)) + } + current := a.sub + c.mu.Unlock() + _ = current.Unsubscribe() + }, nil } // nullable keeps a nil result as JSON null rather than dropping the key: the TypeScript reply always diff --git a/node-tools/internal/bus/refused_test.go b/node-tools/internal/bus/refused_test.go new file mode 100644 index 0000000..026db51 --- /dev/null +++ b/node-tools/internal/bus/refused_test.go @@ -0,0 +1,86 @@ +package bus + +import ( + "encoding/json" + "fmt" + "os" + "strings" + "sync" + "testing" + "time" + + "github.com/nats-io/nats.go" +) + +// novox/hq issue 222: a subscription the bus refused is asked for again. A push sends the machine +// that needs a grant and the bus's machine together; the runtime may subscribe the moment before the +// bus reloads, and a refusal must not leave the module unreachable until some later membership. +func TestARefusedSubscriptionIsAskedForAgain(t *testing.T) { + url := os.Getenv("MESH_TEST_NATS") + if url == "" { + t.Skip("MESH_TEST_NATS unset") + } + c, err := Connect(Credential{URL: url, Module: "alpha", Node: "anchor"}) + if err != nil { + t.Fatal(err) + } + defer c.Close() + var mu sync.Mutex + var said []string + c.Logf = func(f string, a ...any) { mu.Lock(); said = append(said, fmt.Sprintf(f, a...)); mu.Unlock() } + c.RetryAfter = []time.Duration{50 * time.Millisecond, 50 * time.Millisecond} + + subject := "mesh.mod.alpha.tool.ping.anchor" + stop, err := c.answerOn(subject, "", func(json.RawMessage) (any, error) { return "pong", nil }) + if err != nil { + t.Fatal(err) + } + first := c.answering[answeringKey(subject, "")].sub + + // What the bus says when the grant is not there yet. + c.refused(fmt.Errorf("%w: Permissions Violation for Subscription to %q", nats.ErrPermissionViolation, subject)) + deadline := time.Now().Add(2 * time.Second) + for { + c.mu.Lock() + again := c.answering[answeringKey(subject, "")].sub + c.mu.Unlock() + if again != first { + break + } + if time.Now().After(deadline) { + t.Fatal("the refused subscription was not asked for again") + } + time.Sleep(10 * time.Millisecond) + } + asker, err := nats.Connect(url) + if err != nil { + t.Fatal(err) + } + defer asker.Close() + reply, err := asker.Request(subject, []byte("{}"), 2*time.Second) + if err != nil { + t.Fatalf("the subject asked for again does not answer: %v", err) + } + if !strings.Contains(string(reply.Data), "pong") { + t.Fatalf("the subject asked for again answered %s", reply.Data) + } + + // Past its attempts it is given up, and said. + for i := 0; i < 3; i++ { + c.refused(fmt.Errorf("%w: Permissions Violation for Subscription to %q", nats.ErrPermissionViolation, subject)) + time.Sleep(80 * time.Millisecond) + } + mu.Lock() + gaveUp := strings.Contains(strings.Join(said, "\n"), "still refuses "+subject) + mu.Unlock() + if !gaveUp { + t.Errorf("never gave up on a subject that stays refused:\n%s", strings.Join(said, "\n")) + } + + // A subject the runtime stopped answering is not asked for again. + stop() + c.refused(fmt.Errorf("%w: Permissions Violation for Subscription to %q", nats.ErrPermissionViolation, subject)) + if _, still := c.answering[answeringKey(subject, "")]; still { + t.Error("a stopped subject is still tracked") + } +} -- 2.54.0