Bump github.com/nats-io/nats-server/v2 from 2.9.4 to 2.9.17 (#6373)
Bumps [github.com/nats-io/nats-server/v2](https://github.com/nats-io/nats-server) from 2.9.4 to 2.9.17. - [Release notes](https://github.com/nats-io/nats-server/releases) - [Changelog](https://github.com/nats-io/nats-server/blob/main/.goreleaser.yml) - [Commits](https://github.com/nats-io/nats-server/compare/v2.9.4...v2.9.17) --- updated-dependencies: - dependency-name: github.com/nats-io/nats-server/v2 dependency-type: direct:production update-type: version-update:semver-patch ... Signed-off-by: dependabot[bot] <support@github.com> Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com>
This commit is contained in:
co-authored by
dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com>
parent
c0665975b3
commit
1ad0637254
+138
-55
@@ -1,4 +1,4 @@
|
||||
// Copyright 2019-2022 The NATS Authors
|
||||
// Copyright 2019-2023 The NATS Authors
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
@@ -39,29 +39,31 @@ import (
|
||||
"github.com/nats-io/nuid"
|
||||
)
|
||||
|
||||
// Warning when user configures leafnode TLS insecure
|
||||
const leafnodeTLSInsecureWarning = "TLS certificate chain and hostname of solicited leafnodes will not be verified. DO NOT USE IN PRODUCTION!"
|
||||
const (
|
||||
// Warning when user configures leafnode TLS insecure
|
||||
leafnodeTLSInsecureWarning = "TLS certificate chain and hostname of solicited leafnodes will not be verified. DO NOT USE IN PRODUCTION!"
|
||||
|
||||
// When a loop is detected, delay the reconnect of solicited connection.
|
||||
const leafNodeReconnectDelayAfterLoopDetected = 30 * time.Second
|
||||
// When a loop is detected, delay the reconnect of solicited connection.
|
||||
leafNodeReconnectDelayAfterLoopDetected = 30 * time.Second
|
||||
|
||||
// When a server receives a message causing a permission violation, the
|
||||
// connection is closed and it won't attempt to reconnect for that long.
|
||||
const leafNodeReconnectAfterPermViolation = 30 * time.Second
|
||||
// When a server receives a message causing a permission violation, the
|
||||
// connection is closed and it won't attempt to reconnect for that long.
|
||||
leafNodeReconnectAfterPermViolation = 30 * time.Second
|
||||
|
||||
// When we have the same cluster name as the hub.
|
||||
const leafNodeReconnectDelayAfterClusterNameSame = 30 * time.Second
|
||||
// When we have the same cluster name as the hub.
|
||||
leafNodeReconnectDelayAfterClusterNameSame = 30 * time.Second
|
||||
|
||||
// Prefix for loop detection subject
|
||||
const leafNodeLoopDetectionSubjectPrefix = "$LDS."
|
||||
// Prefix for loop detection subject
|
||||
leafNodeLoopDetectionSubjectPrefix = "$LDS."
|
||||
|
||||
// Path added to URL to indicate to WS server that the connection is a
|
||||
// LEAF connection as opposed to a CLIENT.
|
||||
const leafNodeWSPath = "/leafnode"
|
||||
// Path added to URL to indicate to WS server that the connection is a
|
||||
// LEAF connection as opposed to a CLIENT.
|
||||
leafNodeWSPath = "/leafnode"
|
||||
|
||||
// This is the time the server will wait, when receiving a CONNECT,
|
||||
// before closing the connection if the required minimum version is not met.
|
||||
const leafNodeWaitBeforeClose = 5 * time.Second
|
||||
// This is the time the server will wait, when receiving a CONNECT,
|
||||
// before closing the connection if the required minimum version is not met.
|
||||
leafNodeWaitBeforeClose = 5 * time.Second
|
||||
)
|
||||
|
||||
type leaf struct {
|
||||
// We have any auth stuff here for solicited connections.
|
||||
@@ -528,10 +530,45 @@ func (s *Server) connectToRemoteLeafNode(remote *leafNodeCfg, firstConnect bool)
|
||||
// We have a connection here to a remote server.
|
||||
// Go ahead and create our leaf node and return.
|
||||
s.createLeafNode(conn, rURL, remote, nil)
|
||||
|
||||
// Clear any observer states if we had them.
|
||||
s.clearObserverState(remote)
|
||||
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
// This will clear any observer state such that stream or consumer assets on this server can become leaders again.
|
||||
func (s *Server) clearObserverState(remote *leafNodeCfg) {
|
||||
s.mu.RLock()
|
||||
accName := remote.LocalAccount
|
||||
s.mu.RUnlock()
|
||||
|
||||
acc, err := s.LookupAccount(accName)
|
||||
if err != nil {
|
||||
s.Warnf("Error looking up account [%s] checking for JetStream clear observer state on a leafnode", accName)
|
||||
return
|
||||
}
|
||||
|
||||
// Walk all streams looking for any clustered stream, skip otherwise.
|
||||
for _, mset := range acc.streams() {
|
||||
node := mset.raftNode()
|
||||
if node == nil {
|
||||
// Not R>1
|
||||
continue
|
||||
}
|
||||
// Check consumers
|
||||
for _, o := range mset.getConsumers() {
|
||||
if n := o.raftNode(); n != nil {
|
||||
// Ensure we can become a leader again.
|
||||
n.SetObserver(false)
|
||||
}
|
||||
}
|
||||
// Ensure we can not become a leader again.
|
||||
node.SetObserver(false)
|
||||
}
|
||||
}
|
||||
|
||||
// Check to see if we should migrate any assets from this account.
|
||||
func (s *Server) checkJetStreamMigrate(remote *leafNodeCfg) {
|
||||
s.mu.RLock()
|
||||
@@ -544,7 +581,7 @@ func (s *Server) checkJetStreamMigrate(remote *leafNodeCfg) {
|
||||
|
||||
acc, err := s.LookupAccount(accName)
|
||||
if err != nil {
|
||||
s.Debugf("Error looking up account [%s] checking for JetStream migration on a leafnode", accName)
|
||||
s.Warnf("Error looking up account [%s] checking for JetStream migration on a leafnode", accName)
|
||||
return
|
||||
}
|
||||
|
||||
@@ -558,14 +595,20 @@ func (s *Server) checkJetStreamMigrate(remote *leafNodeCfg) {
|
||||
}
|
||||
// Collect any consumers
|
||||
for _, o := range mset.getConsumers() {
|
||||
if n := o.raftNode(); n != nil && n.Leader() {
|
||||
n.StepDown()
|
||||
if n := o.raftNode(); n != nil {
|
||||
if n.Leader() {
|
||||
n.StepDown()
|
||||
}
|
||||
// Ensure we can not become a leader while in this state.
|
||||
n.SetObserver(true)
|
||||
}
|
||||
}
|
||||
// Stepdown if this stream was leader.
|
||||
if node.Leader() {
|
||||
node.StepDown()
|
||||
}
|
||||
// Ensure we can not become a leader while in this state.
|
||||
node.SetObserver(true)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -654,7 +697,7 @@ func (s *Server) startLeafNodeAcceptLoop() {
|
||||
s.leafNodeInfo = info
|
||||
// Possibly override Host/Port and set IP based on Cluster.Advertise
|
||||
if err := s.setLeafNodeInfoHostPortAndIP(); err != nil {
|
||||
s.Fatalf("Error setting leafnode INFO with LeafNode.Advertise value of %s, err=%v", s.opts.LeafNode.Advertise, err)
|
||||
s.Fatalf("Error setting leafnode INFO with LeafNode.Advertise value of %s, err=%v", opts.LeafNode.Advertise, err)
|
||||
l.Close()
|
||||
s.mu.Unlock()
|
||||
return
|
||||
@@ -816,6 +859,7 @@ func (s *Server) removeLeafNodeURL(urlStr string) bool {
|
||||
|
||||
// Server lock is held on entry
|
||||
func (s *Server) generateLeafNodeInfoJSON() {
|
||||
s.leafNodeInfo.Cluster = s.cachedClusterName()
|
||||
s.leafNodeInfo.LeafNodeURLs = s.leafURLsMap.getAsStringSlice()
|
||||
s.leafNodeInfo.WSConnectURLs = s.websocket.connectURLsMap.getAsStringSlice()
|
||||
b, _ := json.Marshal(s.leafNodeInfo)
|
||||
@@ -1339,7 +1383,7 @@ func (s *Server) addLeafNodeConnection(c *client, srvName, clusterName string, c
|
||||
c.mergeDenyPermissionsLocked(both, denyAllJs)
|
||||
// When a remote with a system account is present in a server, unless otherwise disabled, the server will be
|
||||
// started in observer mode. Now that it is clear that this not used, turn the observer mode off.
|
||||
if solicited && meta != nil && meta.isObserver() {
|
||||
if solicited && meta != nil && meta.IsObserver() {
|
||||
meta.setObserver(false, extNotExtended)
|
||||
c.Noticef("Turning JetStream metadata controller Observer Mode off")
|
||||
// Take note that the domain was not extended to avoid this state from startup.
|
||||
@@ -1360,7 +1404,7 @@ func (s *Server) addLeafNodeConnection(c *client, srvName, clusterName string, c
|
||||
myRemoteDomain, srvDecorated())
|
||||
// In an extension use case, pin leadership to server remotes connect to.
|
||||
// Therefore, server with a remote that are not already in observer mode, need to be put into it.
|
||||
if solicited && meta != nil && !meta.isObserver() {
|
||||
if solicited && meta != nil && !meta.IsObserver() {
|
||||
meta.setObserver(true, extExtended)
|
||||
c.Noticef("Turning JetStream metadata controller Observer Mode on - System Account Connected")
|
||||
// Take note that the domain was not extended to avoid this state next startup.
|
||||
@@ -1537,6 +1581,11 @@ func (c *client) processLeafNodeConnect(s *Server, arg []byte, lang string) erro
|
||||
|
||||
c.mu.Unlock()
|
||||
|
||||
// Register the cluster, even if empty, as long as we are acting as a hub.
|
||||
if !proto.Hub {
|
||||
c.acc.registerLeafNodeCluster(proto.Cluster)
|
||||
}
|
||||
|
||||
// Add in the leafnode here since we passed through auth at this point.
|
||||
s.addLeafNodeConnection(c, proto.Name, proto.Cluster, true)
|
||||
|
||||
@@ -1590,11 +1639,11 @@ func (s *Server) initLeafNodeSmapAndSendSubs(c *client) {
|
||||
return
|
||||
}
|
||||
// Collect all account subs here.
|
||||
_subs := [32]*subscription{}
|
||||
_subs := [1024]*subscription{}
|
||||
subs := _subs[:0]
|
||||
ims := []string{}
|
||||
|
||||
acc.mu.Lock()
|
||||
acc.mu.RLock()
|
||||
accName := acc.Name
|
||||
accNTag := acc.nameTag
|
||||
|
||||
@@ -1633,11 +1682,15 @@ func (s *Server) initLeafNodeSmapAndSendSubs(c *client) {
|
||||
|
||||
// Create a unique subject that will be used for loop detection.
|
||||
lds := acc.lds
|
||||
acc.mu.RUnlock()
|
||||
|
||||
// Check if we have to create the LDS.
|
||||
if lds == _EMPTY_ {
|
||||
lds = leafNodeLoopDetectionSubjectPrefix + nuid.Next()
|
||||
acc.mu.Lock()
|
||||
acc.lds = lds
|
||||
acc.mu.Unlock()
|
||||
}
|
||||
acc.mu.Unlock()
|
||||
|
||||
// Now check for gateway interest. Leafnodes will put this into
|
||||
// the proper mode to propagate, but they are not held in the account.
|
||||
@@ -1741,48 +1794,76 @@ func (s *Server) updateInterestForAccountOnGateway(accName string, sub *subscrip
|
||||
s.Debugf("No or bad account for %q, failed to update interest from gateway", accName)
|
||||
return
|
||||
}
|
||||
s.updateLeafNodes(acc, sub, delta)
|
||||
acc.updateLeafNodes(sub, delta)
|
||||
}
|
||||
|
||||
// updateLeafNodes will make sure to update the smap for the subscription. Will
|
||||
// also forward to all leaf nodes as needed.
|
||||
func (s *Server) updateLeafNodes(acc *Account, sub *subscription, delta int32) {
|
||||
// updateLeafNodes will make sure to update the account smap for the subscription.
|
||||
// Will also forward to all leaf nodes as needed.
|
||||
func (acc *Account) updateLeafNodes(sub *subscription, delta int32) {
|
||||
if acc == nil || sub == nil {
|
||||
return
|
||||
}
|
||||
|
||||
_l := [32]*client{}
|
||||
leafs := _l[:0]
|
||||
// We will do checks for no leafnodes and same cluster here inline and under the
|
||||
// general account read lock.
|
||||
// If we feel we need to update the leafnodes we will do that out of line to avoid
|
||||
// blocking routes or GWs.
|
||||
|
||||
// Grab all leaf nodes. Ignore a leafnode if sub's client is a leafnode and matches.
|
||||
acc.mu.RLock()
|
||||
for _, ln := range acc.lleafs {
|
||||
if ln != sub.client {
|
||||
leafs = append(leafs, ln)
|
||||
}
|
||||
// First check if we even have leafnodes here.
|
||||
if acc.nleafs == 0 {
|
||||
acc.mu.RUnlock()
|
||||
return
|
||||
}
|
||||
|
||||
// Is this a loop detection subject.
|
||||
isLDS := bytes.HasPrefix(sub.subject, []byte(leafNodeLoopDetectionSubjectPrefix))
|
||||
|
||||
// Capture the cluster even if its empty.
|
||||
cluster := _EMPTY_
|
||||
if sub.origin != nil {
|
||||
cluster = string(sub.origin)
|
||||
}
|
||||
|
||||
// If we have an isolated cluster we can return early, as long as it is not a loop detection subject.
|
||||
// Empty clusters will return false for the check.
|
||||
if !isLDS && acc.isLeafNodeClusterIsolated(cluster) {
|
||||
acc.mu.RUnlock()
|
||||
return
|
||||
}
|
||||
|
||||
// We can release the general account lock.
|
||||
acc.mu.RUnlock()
|
||||
|
||||
for _, ln := range leafs {
|
||||
// Check to make sure this sub does not have an origin cluster than matches the leafnode.
|
||||
ln.mu.Lock()
|
||||
skip := (sub.origin != nil && string(sub.origin) == ln.remoteCluster()) || !ln.canSubscribe(string(sub.subject))
|
||||
ln.mu.Unlock()
|
||||
if skip {
|
||||
// We can hold the list lock here to avoid having to copy a large slice.
|
||||
acc.lmu.RLock()
|
||||
defer acc.lmu.RUnlock()
|
||||
|
||||
// Do this once.
|
||||
subject := string(sub.subject)
|
||||
|
||||
// Walk the connected leafnodes.
|
||||
for _, ln := range acc.lleafs {
|
||||
if ln == sub.client {
|
||||
continue
|
||||
}
|
||||
ln.updateSmap(sub, delta)
|
||||
// Check to make sure this sub does not have an origin cluster that matches the leafnode.
|
||||
ln.mu.Lock()
|
||||
skip := (cluster != _EMPTY_ && cluster == ln.remoteCluster()) || (delta > 0 && !ln.canSubscribe(subject))
|
||||
// If skipped, make sure that we still let go the "$LDS." subscription that allows
|
||||
// the detection of a loop.
|
||||
if isLDS || !skip {
|
||||
ln.updateSmap(sub, delta)
|
||||
}
|
||||
ln.mu.Unlock()
|
||||
}
|
||||
}
|
||||
|
||||
// This will make an update to our internal smap and determine if we should send out
|
||||
// an interest update to the remote side.
|
||||
// Lock should be held.
|
||||
func (c *client) updateSmap(sub *subscription, delta int32) {
|
||||
key := keyFromSub(sub)
|
||||
|
||||
c.mu.Lock()
|
||||
if c.leaf.smap == nil {
|
||||
c.mu.Unlock()
|
||||
return
|
||||
}
|
||||
|
||||
@@ -1790,7 +1871,6 @@ func (c *client) updateSmap(sub *subscription, delta int32) {
|
||||
skind := sub.client.kind
|
||||
updateClient := skind == CLIENT || skind == SYSTEM || skind == JETSTREAM || skind == ACCOUNT
|
||||
if c.isSpokeLeafNode() && !(updateClient || (skind == LEAF && !sub.client.isSpokeLeafNode())) {
|
||||
c.mu.Unlock()
|
||||
return
|
||||
}
|
||||
|
||||
@@ -1803,12 +1883,16 @@ func (c *client) updateSmap(sub *subscription, delta int32) {
|
||||
c.leaf.tsubt.Stop()
|
||||
c.leaf.tsubt = nil
|
||||
}
|
||||
c.mu.Unlock()
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
n := c.leaf.smap[key]
|
||||
key := keyFromSub(sub)
|
||||
n, ok := c.leaf.smap[key]
|
||||
if delta < 0 && !ok {
|
||||
return
|
||||
}
|
||||
|
||||
// We will update if its a queue, if count is zero (or negative), or we were 0 and are N > 0.
|
||||
update := sub.queue != nil || n == 0 || n+delta <= 0
|
||||
n += delta
|
||||
@@ -1820,7 +1904,6 @@ func (c *client) updateSmap(sub *subscription, delta int32) {
|
||||
if update {
|
||||
c.sendLeafNodeSubUpdate(key, n)
|
||||
}
|
||||
c.mu.Unlock()
|
||||
}
|
||||
|
||||
// Used to force add subjects to the subject map.
|
||||
@@ -2034,7 +2117,7 @@ func (c *client) processLeafSub(argo []byte) (err error) {
|
||||
}
|
||||
// Now check on leafnode updates for other leaf nodes. We understand solicited
|
||||
// and non-solicited state in this call so we will do the right thing.
|
||||
srv.updateLeafNodes(acc, sub, delta)
|
||||
acc.updateLeafNodes(sub, delta)
|
||||
|
||||
return nil
|
||||
}
|
||||
@@ -2091,7 +2174,7 @@ func (c *client) processLeafUnsub(arg []byte) error {
|
||||
}
|
||||
}
|
||||
// Now check on leafnode updates for other leaf nodes.
|
||||
srv.updateLeafNodes(acc, sub, -1)
|
||||
acc.updateLeafNodes(sub, -1)
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user