Ask the bus again for a subscription it refused (hq issue 222) #47
@@ -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
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user