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") + } +}