• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

nats-io / nats-server / 30145931992

24 Jul 2026 09:07AM UTC coverage: 77.402% (-0.09%) from 77.489%
30145931992

push

github

web-flow
[FIXED] MQTT: QoS2 message released on a resumed session lost its QoS and PI (#8414)

This came up while testing with
https://github.com/ConnectEverything/mqtt-test/pull/4

Two root causes. The delivery callbacks and the retain handling read the
publishing connection's c.mqtt.pp — the last PUBLISH parsed — which is
stale for a PUBREL- or will-initiated delivery, so a QoS2 message
released after a reconnect was delivered as QoS0 with no packet
identifier. And the staged copy in $MQTT_qos2in never persisted the
retain flag, so it could not be recovered on another connection at all.
Same-connection flows worked only while the scratch happened to still
hold the right packet. Because of the interplay, the two must be fixed
together.

The delivered message now becomes the client's current publish for the
duration of the broadcast (swapped into the readLoop-owned c.mqtt.pp,
restored by the same defer as the pubargs); wills build their own
mqttPublish and use the common path instead of writing into the scratch
by hand. The staged copy
persists the retain flag in a second, flags byte of the Nmqtt-Pub header
value, restored when the PUBREL rebuilds the publish; older servers read
only the first (QoS) byte, and messages without a flags byte read as not
retained.

Retained QoS2 messages now survive an in-flight reconnect: a PUBLISH
staged on one connection and released by the PUBREL on the next stores
the retained message, where before the retain bit died with the first
connection's parser state. This matches Spec [MQTT-4.3.3] Method A as
implemented by mosquitto, which likewise persists the retain flag with
the staged message and applies it at release, together with the
delivery.

Added regression tests: a PUBREL on a resumed session delivers at the
subscription's QoS exactly once, a retransmitted PUBREL is acknowledged
without redelivering, and the retain flag survives the staging round
trip.

73227 of 94606 relevant lines covered (77.4%)

461998.16 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

87.69
/server/leafnode.go
1
// Copyright 2019-2026 The NATS Authors
2
// Licensed under the Apache License, Version 2.0 (the "License");
3
// you may not use this file except in compliance with the License.
4
// You may obtain a copy of the License at
5
//
6
// http://www.apache.org/licenses/LICENSE-2.0
7
//
8
// Unless required by applicable law or agreed to in writing, software
9
// distributed under the License is distributed on an "AS IS" BASIS,
10
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
11
// See the License for the specific language governing permissions and
12
// limitations under the License.
13

14
package server
15

16
import (
17
        "bufio"
18
        "bytes"
19
        "crypto/tls"
20
        "encoding/base64"
21
        "encoding/json"
22
        "fmt"
23
        "io"
24
        "math/rand"
25
        "net"
26
        "net/http"
27
        "net/url"
28
        "os"
29
        "path"
30
        "regexp"
31
        "runtime"
32
        "strconv"
33
        "strings"
34
        "sync"
35
        "sync/atomic"
36
        "time"
37

38
        "github.com/klauspost/compress/s2"
39
        "github.com/nats-io/jwt/v2"
40
        "github.com/nats-io/nkeys"
41
        "github.com/nats-io/nuid"
42
)
43

44
const (
45
        // Warning when user configures leafnode TLS insecure
46
        leafnodeTLSInsecureWarning = "TLS certificate chain and hostname of solicited leafnodes will not be verified. DO NOT USE IN PRODUCTION!"
47

48
        // When a loop is detected, delay the reconnect of solicited connection.
49
        leafNodeReconnectDelayAfterLoopDetected = 30 * time.Second
50

51
        // When a server receives a message causing a permission violation, the
52
        // connection is closed and it won't attempt to reconnect for that long.
53
        leafNodeReconnectAfterPermViolation = 30 * time.Second
54

55
        // When we have the same cluster name as the hub.
56
        leafNodeReconnectDelayAfterClusterNameSame = 30 * time.Second
57

58
        // Prefix for loop detection subject
59
        leafNodeLoopDetectionSubjectPrefix = "$LDS."
60

61
        // Path added to URL to indicate to WS server that the connection is a
62
        // LEAF connection as opposed to a CLIENT.
63
        leafNodeWSPath = "/leafnode"
64

65
        // When a soliciting leafnode is rejected because it does not meet the
66
        // configured minimum version, delay the next reconnect attempt by this long.
67
        leafNodeMinVersionReconnectDelay = 5 * time.Second
68
)
69

70
type leaf struct {
71
        // We have any auth stuff here for solicited connections.
72
        remote *leafNodeCfg
73
        // isSpoke tells us what role we are playing.
74
        // Used when we receive a connection but otherside tells us they are a hub.
75
        isSpoke bool
76
        // remoteCluster is when we are a hub but the spoke leafnode is part of a cluster.
77
        remoteCluster string
78
        // remoteServer holds onto the remote server's name or ID.
79
        remoteServer string
80
        // domain name of remote server
81
        remoteDomain string
82
        // account name of remote server
83
        remoteAccName string
84
        // Whether or not we want to propagate east-west interest from other LNs.
85
        isolated bool
86
        // Used to suppress sub and unsub interest. Same as routes but our audience
87
        // here is tied to this leaf node. This will hold all subscriptions except this
88
        // leaf nodes. This represents all the interest we want to send to the other side.
89
        smap map[string]int32
90
        // This map will contain all the subscriptions that have been added to the smap
91
        // during initLeafNodeSmapAndSendSubs. It is short lived and is there to avoid
92
        // race between processing of a sub where sub is added to account sublist but
93
        // updateSmap has not be called on that "thread", while in the LN readloop,
94
        // when processing CONNECT, initLeafNodeSmapAndSendSubs is invoked and add
95
        // this subscription to smap. When processing of the sub then calls updateSmap,
96
        // we would add it a second time in the smap causing later unsub to suppress the LS-.
97
        tsub  map[*subscription]struct{}
98
        tsubt *time.Timer
99
        // Selected compression mode, which may be different from the server configured mode.
100
        compression string
101
        // This is for GW map replies.
102
        gwSub *subscription
103
}
104

105
// Used for remote (solicited) leafnodes.
106
type leafNodeCfg struct {
107
        sync.RWMutex
108
        *RemoteLeafOpts
109
        urls           []*url.URL
110
        curURL         *url.URL
111
        tlsName        string
112
        username       string
113
        password       string
114
        perms          *Permissions
115
        connDelay      time.Duration // Delay before a connect, could be used while detecting loop condition, etc..
116
        jsMigrateTimer *time.Timer
117
        quitCh         chan struct{}
118
        removed        bool
119
        connInProgress bool
120
}
121

122
// Check to see if this is a solicited leafnode. We do special processing for solicited.
123
func (c *client) isSolicitedLeafNode() bool {
1,779✔
124
        return c.kind == LEAF && c.leaf != nil && c.leaf.remote != nil
1,779✔
125
}
1,779✔
126

127
// Returns true if this is a solicited leafnode and is not configured to be treated as a hub or a receiving
128
// connection leafnode where the otherside has declared itself to be the hub.
129
func (c *client) isSpokeLeafNode() bool {
9,322,330✔
130
        return c.kind == LEAF && c.leaf != nil && c.leaf.isSpoke
9,322,330✔
131
}
9,322,330✔
132

133
func (c *client) isHubLeafNode() bool {
15,750✔
134
        return c.kind == LEAF && c.leaf != nil && !c.leaf.isSpoke
15,750✔
135
}
15,750✔
136

137
func (c *client) isIsolatedLeafNode() bool {
8,981✔
138
        // TODO(nat): In future we may want to pass in and consider an isolation
8,981✔
139
        // group name here, which the hub and/or leaf could provide, so that we
8,981✔
140
        // can isolate away certain LNs but not others on an opt-in basis. For
8,981✔
141
        // now we will just isolate all LN interest until then.
8,981✔
142
        return c.kind == LEAF && c.leaf != nil && c.leaf.isolated
8,981✔
143
}
8,981✔
144

145
// This will spin up go routines to solicit the remote leaf node connections.
146
func (s *Server) solicitLeafNodeRemotes(remotes []*RemoteLeafOpts) {
469✔
147
        sysAccName := _EMPTY_
469✔
148
        sAcc := s.SystemAccount()
469✔
149
        if sAcc != nil {
916✔
150
                sysAccName = sAcc.Name
447✔
151
        }
447✔
152
        addRemote := func(r *RemoteLeafOpts, isSysAccRemote bool) *leafNodeCfg {
1,049✔
153
                s.mu.Lock()
580✔
154
                remote := newLeafNodeCfg(r)
580✔
155
                creds := remote.Credentials
580✔
156
                accName := remote.LocalAccount
580✔
157
                if s.leafRemoteCfgs == nil {
1,048✔
158
                        s.leafRemoteCfgs = make(map[*leafNodeCfg]struct{})
468✔
159
                }
468✔
160
                s.leafRemoteCfgs[remote] = struct{}{}
580✔
161
                // Print notice if
580✔
162
                if isSysAccRemote {
666✔
163
                        if len(remote.DenyExports) > 0 {
87✔
164
                                s.Noticef("Remote for System Account uses restricted export permissions")
1✔
165
                        }
1✔
166
                        if len(remote.DenyImports) > 0 {
87✔
167
                                s.Noticef("Remote for System Account uses restricted import permissions")
1✔
168
                        }
1✔
169
                }
170
                s.mu.Unlock()
580✔
171
                if creds != _EMPTY_ {
620✔
172
                        contents, err := os.ReadFile(creds)
40✔
173
                        defer wipeSlice(contents)
40✔
174
                        if err != nil {
40✔
175
                                s.Errorf("Error reading LeafNode Remote Credentials file %q: %v", creds, err)
×
176
                        } else if items := credsRe.FindAllSubmatch(contents, -1); len(items) < 2 {
40✔
177
                                s.Errorf("LeafNode Remote Credentials file %q malformed", creds)
×
178
                        } else if _, err := nkeys.FromSeed(items[1][1]); err != nil {
40✔
179
                                s.Errorf("LeafNode Remote Credentials file %q has malformed seed", creds)
×
180
                        } else if uc, err := jwt.DecodeUserClaims(string(items[0][1])); err != nil {
40✔
181
                                s.Errorf("LeafNode Remote Credentials file %q has malformed user jwt", creds)
×
182
                        } else if isSysAccRemote {
44✔
183
                                if !uc.Permissions.Pub.Empty() || !uc.Permissions.Sub.Empty() || uc.Permissions.Resp != nil {
5✔
184
                                        s.Noticef("LeafNode Remote for System Account uses credentials file %q with restricted permissions", creds)
1✔
185
                                }
1✔
186
                        } else {
36✔
187
                                if !uc.Permissions.Pub.Empty() || !uc.Permissions.Sub.Empty() || uc.Permissions.Resp != nil {
40✔
188
                                        s.Noticef("LeafNode Remote for Account %s uses credentials file %q with restricted permissions", accName, creds)
4✔
189
                                }
4✔
190
                        }
191
                }
192
                return remote
580✔
193
        }
194
        for _, r := range remotes {
1,049✔
195
                // We need to call this, even if the leaf is disabled. This is so that
580✔
196
                // the number of internal configuration matches the options' remote leaf
580✔
197
                // configuration required for configuration reload.
580✔
198
                remote := addRemote(r, r.LocalAccount == sysAccName)
580✔
199
                if !r.Disabled {
1,160✔
200
                        s.connectToRemoteLeafNodeAsynchronously(remote, true)
580✔
201
                }
580✔
202
        }
203
}
204

205
// Ensure that leafnode is properly configured.
206
func validateLeafNode(o *Options) error {
7,753✔
207
        if err := validateLeafNodeAuthOptions(o); err != nil {
7,755✔
208
                return err
2✔
209
        }
2✔
210

211
        if len(o.LeafNode.Remotes) > 0 {
8,251✔
212
                names := make(map[string]struct{})
500✔
213
                // Check for duplicate remotes, also, users can bind to any local account,
500✔
214
                // if its empty we will assume the $G account.
500✔
215
                for _, r := range o.LeafNode.Remotes {
1,117✔
216
                        if r.LocalAccount == _EMPTY_ {
995✔
217
                                r.LocalAccount = globalAccountName
378✔
218
                        }
378✔
219
                        rn := r.name()
617✔
220
                        if _, dup := names[rn]; dup {
620✔
221
                                return fmt.Errorf("duplicate remote %s", r.safeName())
3✔
222
                        }
3✔
223
                        names[rn] = struct{}{}
614✔
224
                }
225
        }
226

227
        // In local config mode, check that leafnode configuration refers to accounts that exist.
228
        if len(o.TrustedOperators) == 0 {
15,215✔
229
                accNames := map[string]struct{}{}
7,467✔
230
                for _, a := range o.Accounts {
15,973✔
231
                        accNames[a.Name] = struct{}{}
8,506✔
232
                }
8,506✔
233
                // global account is always created
234
                accNames[DEFAULT_GLOBAL_ACCOUNT] = struct{}{}
7,467✔
235
                // in the context of leaf nodes, empty account means global account
7,467✔
236
                accNames[_EMPTY_] = struct{}{}
7,467✔
237
                // system account either exists or, if not disabled, will be created
7,467✔
238
                if o.SystemAccount == _EMPTY_ && !o.NoSystemAccount {
13,538✔
239
                        accNames[DEFAULT_SYSTEM_ACCOUNT] = struct{}{}
6,071✔
240
                }
6,071✔
241
                checkAccountExists := func(accName string, cfgType string) error {
15,548✔
242
                        if _, ok := accNames[accName]; !ok {
8,083✔
243
                                return fmt.Errorf("cannot find local account %q specified in leafnode %s", accName, cfgType)
2✔
244
                        }
2✔
245
                        return nil
8,079✔
246
                }
247
                if err := checkAccountExists(o.LeafNode.Account, "authorization"); err != nil {
7,468✔
248
                        return err
1✔
249
                }
1✔
250
                for _, lu := range o.LeafNode.Users {
7,483✔
251
                        if lu.Account == nil { // means global account
27✔
252
                                continue
10✔
253
                        }
254
                        if err := checkAccountExists(lu.Account.Name, "authorization"); err != nil {
7✔
255
                                return err
×
256
                        }
×
257
                }
258
                for _, r := range o.LeafNode.Remotes {
8,073✔
259
                        if err := checkAccountExists(r.LocalAccount, "remote"); err != nil {
608✔
260
                                return err
1✔
261
                        }
1✔
262
                }
263
        } else {
281✔
264
                if len(o.LeafNode.Users) != 0 {
282✔
265
                        return fmt.Errorf("operator mode does not allow specifying users in leafnode config")
1✔
266
                }
1✔
267
                for _, r := range o.LeafNode.Remotes {
281✔
268
                        if !nkeys.IsValidPublicAccountKey(r.LocalAccount) {
2✔
269
                                return fmt.Errorf(
1✔
270
                                        "operator mode requires account nkeys in remotes. " +
1✔
271
                                                "Please add an `account` key to each remote in your `leafnodes` section, to assign it to an account. " +
1✔
272
                                                "Each account value should be a 56 character public key, starting with the letter 'A'")
1✔
273
                        }
1✔
274
                }
275
                if o.LeafNode.Port != 0 && o.LeafNode.Account != "" && !nkeys.IsValidPublicAccountKey(o.LeafNode.Account) {
280✔
276
                        return fmt.Errorf("operator mode and non account nkeys are incompatible")
1✔
277
                }
1✔
278
        }
279

280
        // Validate compression settings
281
        if o.LeafNode.Compression.Mode != _EMPTY_ {
12,334✔
282
                if err := validateAndNormalizeCompressionOption(&o.LeafNode.Compression, CompressionS2Auto); err != nil {
4,596✔
283
                        return err
5✔
284
                }
5✔
285
        }
286

287
        // If a remote has a websocket scheme, all need to have it.
288
        for _, rcfg := range o.LeafNode.Remotes {
8,344✔
289
                // Validate proxy configuration
606✔
290
                if _, err := validateLeafNodeProxyOptions(rcfg); err != nil {
612✔
291
                        return err
6✔
292
                }
6✔
293

294
                if len(rcfg.URLs) >= 2 {
763✔
295
                        firstIsWS, ok := isWSURL(rcfg.URLs[0]), true
163✔
296
                        for i := 1; i < len(rcfg.URLs); i++ {
478✔
297
                                u := rcfg.URLs[i]
315✔
298
                                if isWS := isWSURL(u); isWS && !firstIsWS || !isWS && firstIsWS {
322✔
299
                                        ok = false
7✔
300
                                        break
7✔
301
                                }
302
                        }
303
                        if !ok {
170✔
304
                                return fmt.Errorf("remote leaf node configuration cannot have a mix of websocket and non-websocket urls: %q", redactURLList(rcfg.URLs))
7✔
305
                        }
7✔
306
                }
307
                if !wsAllowedFIPS() {
593✔
308
                        for _, u := range rcfg.URLs {
×
309
                                if isWSURL(u) {
×
310
                                        return fmt.Errorf("remote leaf node URL %q cannot be used in FIPS-140 mode when built with this Go version, use Go 1.26 or later", redactURLString(u.String()))
×
311
                                }
×
312
                        }
313
                }
314
                // Validate compression settings
315
                if rcfg.Compression.Mode != _EMPTY_ {
1,180✔
316
                        if err := validateAndNormalizeCompressionOption(&rcfg.Compression, CompressionS2Auto); err != nil {
592✔
317
                                return err
5✔
318
                        }
5✔
319
                }
320
        }
321

322
        if o.LeafNode.Port == 0 {
11,438✔
323
                return nil
3,718✔
324
        }
3,718✔
325

326
        // If MinVersion is defined, check that it is valid.
327
        if mv := o.LeafNode.MinVersion; mv != _EMPTY_ {
4,006✔
328
                if err := checkLeafMinVersionConfig(mv); err != nil {
6✔
329
                        return err
2✔
330
                }
2✔
331
        }
332

333
        // The checks below will be done only when detecting that we are configured
334
        // with gateways. So if an option validation needs to be done regardless,
335
        // it MUST be done before this point!
336

337
        if o.Gateway.Name == _EMPTY_ && o.Gateway.Port == 0 {
7,362✔
338
                return nil
3,362✔
339
        }
3,362✔
340
        // If we are here we have both leaf nodes and gateways defined, make sure there
341
        // is a system account defined.
342
        if o.SystemAccount == _EMPTY_ {
639✔
343
                return fmt.Errorf("leaf nodes and gateways (both being defined) require a system account to also be configured")
1✔
344
        }
1✔
345
        if err := validatePinnedCerts(o.LeafNode.TLSPinnedCerts); err != nil {
637✔
346
                return fmt.Errorf("leafnode: %v", err)
×
347
        }
×
348
        return nil
637✔
349
}
350

351
func checkLeafMinVersionConfig(mv string) error {
8✔
352
        if ok, err := versionAtLeastCheckError(mv, 2, 8, 0); !ok || err != nil {
12✔
353
                if err != nil {
6✔
354
                        return fmt.Errorf("invalid leafnode's minimum version: %v", err)
2✔
355
                } else {
4✔
356
                        return fmt.Errorf("the minimum version should be at least 2.8.0")
2✔
357
                }
2✔
358
        }
359
        return nil
4✔
360
}
361

362
// Used to validate user names in LeafNode configuration.
363
// - rejects mix of single and multiple users.
364
// - rejects duplicate user names.
365
func validateLeafNodeAuthOptions(o *Options) error {
7,799✔
366
        if len(o.LeafNode.Users) == 0 {
15,572✔
367
                return nil
7,773✔
368
        }
7,773✔
369
        if o.LeafNode.Username != _EMPTY_ {
28✔
370
                return fmt.Errorf("can not have a single user/pass and a users array")
2✔
371
        }
2✔
372
        if o.LeafNode.Nkey != _EMPTY_ {
24✔
373
                return fmt.Errorf("can not have a single nkey and a users array")
×
374
        }
×
375
        users := map[string]struct{}{}
24✔
376
        for _, u := range o.LeafNode.Users {
62✔
377
                if _, exists := users[u.Username]; exists {
40✔
378
                        return fmt.Errorf("duplicate user %q detected in leafnode authorization", u.Username)
2✔
379
                }
2✔
380
                users[u.Username] = struct{}{}
36✔
381
        }
382
        return nil
22✔
383
}
384

385
func validateLeafNodeProxyOptions(remote *RemoteLeafOpts) ([]string, error) {
1,064✔
386
        var warnings []string
1,064✔
387

1,064✔
388
        if remote.Proxy.URL == _EMPTY_ {
2,102✔
389
                return warnings, nil
1,038✔
390
        }
1,038✔
391

392
        proxyURL, err := url.Parse(remote.Proxy.URL)
26✔
393
        if err != nil {
27✔
394
                return warnings, fmt.Errorf("invalid proxy URL: %v", err)
1✔
395
        }
1✔
396

397
        if proxyURL.Scheme != "http" && proxyURL.Scheme != "https" {
27✔
398
                return warnings, fmt.Errorf("proxy URL scheme must be http or https, got: %s", proxyURL.Scheme)
2✔
399
        }
2✔
400

401
        if proxyURL.Host == _EMPTY_ {
25✔
402
                return warnings, fmt.Errorf("proxy URL must specify a host")
2✔
403
        }
2✔
404

405
        if remote.Proxy.Timeout < 0 {
22✔
406
                return warnings, fmt.Errorf("proxy timeout must be >= 0")
1✔
407
        }
1✔
408

409
        if (remote.Proxy.Username == _EMPTY_) != (remote.Proxy.Password == _EMPTY_) {
24✔
410
                return warnings, fmt.Errorf("proxy username and password must both be specified or both be empty")
4✔
411
        }
4✔
412

413
        if len(remote.URLs) > 0 {
32✔
414
                hasWebSocketURL := false
16✔
415
                hasNonWebSocketURL := false
16✔
416

16✔
417
                for _, remoteURL := range remote.URLs {
33✔
418
                        if remoteURL.Scheme == wsSchemePrefix || remoteURL.Scheme == wsSchemePrefixTLS {
30✔
419
                                hasWebSocketURL = true
13✔
420
                                if (remoteURL.Scheme == wsSchemePrefixTLS) &&
13✔
421
                                        remote.TLSConfig == nil && !remote.TLS {
14✔
422
                                        return warnings, fmt.Errorf("proxy is configured but remote URL %s requires TLS and no TLS configuration is provided. When using proxy with TLS endpoints, ensure TLS is properly configured for the leafnode remote", remoteURL.String())
1✔
423
                                }
1✔
424
                        } else {
4✔
425
                                hasNonWebSocketURL = true
4✔
426
                        }
4✔
427
                }
428

429
                if !hasWebSocketURL {
18✔
430
                        warnings = append(warnings, "proxy configuration will be ignored: proxy settings only apply to WebSocket connections (ws:// or wss://), but all configured URLs use TCP connections (nats://)")
3✔
431
                } else if hasNonWebSocketURL {
16✔
432
                        warnings = append(warnings, "proxy configuration will only be used for WebSocket URLs: proxy settings do not apply to TCP connections (nats://)")
1✔
433
                }
1✔
434
        }
435

436
        return warnings, nil
15✔
437
}
438

439
// Wait for the configured reconnect interval before attempting to connect
440
// again to the remote leafnode.
441
func (s *Server) reConnectToRemoteLeafNode(remote *leafNodeCfg) {
239✔
442
        clearInProgress := true
239✔
443
        defer func() {
477✔
444
                s.grWG.Done()
238✔
445
                if clearInProgress {
301✔
446
                        remote.setConnectInProgress(false)
63✔
447
                }
63✔
448
        }()
449
        delay := s.getOpts().LeafNode.ReconnectInterval
239✔
450
        select {
239✔
451
        case <-time.After(delay):
183✔
452
        case <-remote.quitCh:
×
453
                return
×
454
        case <-s.quitCh:
56✔
455
                return
56✔
456
        }
457
        clearInProgress = !connectToRemoteLeafNode(s, remote, false)
183✔
458
}
459

460
// Creates a leafNodeCfg object that wraps the RemoteLeafOpts.
461
func newLeafNodeCfg(remote *RemoteLeafOpts) *leafNodeCfg {
580✔
462
        cfg := &leafNodeCfg{
580✔
463
                RemoteLeafOpts: remote,
580✔
464
                urls:           make([]*url.URL, 0, len(remote.URLs)),
580✔
465
                quitCh:         make(chan struct{}, 1),
580✔
466
        }
580✔
467
        if len(remote.DenyExports) > 0 || len(remote.DenyImports) > 0 {
585✔
468
                perms := &Permissions{}
5✔
469
                if len(remote.DenyExports) > 0 {
10✔
470
                        perms.Publish = &SubjectPermission{Deny: remote.DenyExports}
5✔
471
                }
5✔
472
                if len(remote.DenyImports) > 0 {
9✔
473
                        perms.Subscribe = &SubjectPermission{Deny: remote.DenyImports}
4✔
474
                }
4✔
475
                cfg.perms = perms
5✔
476
        }
477
        // Start with the one that is configured. We will add to this
478
        // array when receiving async leafnode INFOs.
479
        cfg.urls = append(cfg.urls, cfg.URLs...)
580✔
480
        // If allowed to randomize, do it on our copy of URLs
580✔
481
        if !remote.NoRandomize {
1,159✔
482
                rand.Shuffle(len(cfg.urls), func(i, j int) {
871✔
483
                        cfg.urls[i], cfg.urls[j] = cfg.urls[j], cfg.urls[i]
292✔
484
                })
292✔
485
        }
486
        // If we are TLS make sure we save off a proper servername if possible.
487
        // Do same for user/password since we may need them to connect to
488
        // a bare URL that we get from INFO protocol.
489
        for _, u := range cfg.urls {
1,467✔
490
                cfg.saveTLSHostname(u)
887✔
491
                cfg.saveUserPassword(u)
887✔
492
                // If the url(s) have the "wss://" scheme, and we don't have a TLS
887✔
493
                // config, mark that we should be using TLS anyway.
887✔
494
                if !cfg.TLS && isWSSURL(u) {
888✔
495
                        cfg.TLS = true
1✔
496
                }
1✔
497
        }
498
        return cfg
580✔
499
}
500

501
// Notifies the quit channel without blocking.
502
// No lock is needed to invoke this function.
503
func (cfg *leafNodeCfg) notifyQuitChannel() {
2✔
504
        select {
2✔
505
        case cfg.quitCh <- struct{}{}:
2✔
506
        default:
×
507
        }
508
}
509

510
// Sets the connect-in-progress status for this remote leaf configuration.
511
func (cfg *leafNodeCfg) setConnectInProgress(inProgress bool) {
1,907✔
512
        cfg.Lock()
1,907✔
513
        defer cfg.Unlock()
1,907✔
514
        // In both cases we want to drain the "quit" channel.
1,907✔
515
        select {
1,907✔
516
        case <-cfg.quitCh:
×
517
        default:
1,907✔
518
        }
519
        cfg.connInProgress = inProgress
1,907✔
520
}
521

522
// Returns `true` if this remote is in the middle of a connect, `false` otherwise.
523
func (cfg *leafNodeCfg) isConnectInProgress() bool {
2✔
524
        cfg.RLock()
2✔
525
        defer cfg.RUnlock()
2✔
526
        return cfg.connInProgress
2✔
527
}
2✔
528

529
// Mark that this remote is being removed from the configuration.
530
func (cfg *leafNodeCfg) markAsRemoved() {
2✔
531
        cfg.Lock()
2✔
532
        defer cfg.Unlock()
2✔
533
        // This function should be invoked only once, but protect.
2✔
534
        if cfg.removed {
2✔
535
                return
×
536
        }
×
537
        cfg.removed = true
2✔
538
        cfg.notifyQuitChannel()
2✔
539
}
540

541
// Returns false if it has been disabled or removed.
542
func (cfg *leafNodeCfg) stillValid() bool {
5,088✔
543
        cfg.RLock()
5,088✔
544
        defer cfg.RUnlock()
5,088✔
545
        return !cfg.Disabled && !cfg.removed
5,088✔
546
}
5,088✔
547

548
// Will pick an URL from the list of available URLs.
549
func (cfg *leafNodeCfg) pickNextURL() *url.URL {
3,145✔
550
        cfg.Lock()
3,145✔
551
        defer cfg.Unlock()
3,145✔
552
        // If the current URL is the first in the list and we have more than
3,145✔
553
        // one URL, then move that one to end of the list.
3,145✔
554
        if cfg.curURL != nil && len(cfg.urls) > 1 && urlsAreEqual(cfg.curURL, cfg.urls[0]) {
5,519✔
555
                first := cfg.urls[0]
2,374✔
556
                copy(cfg.urls, cfg.urls[1:])
2,374✔
557
                cfg.urls[len(cfg.urls)-1] = first
2,374✔
558
        }
2,374✔
559
        cfg.curURL = cfg.urls[0]
3,145✔
560
        return cfg.curURL
3,145✔
561
}
562

563
// Returns the current URL
564
func (cfg *leafNodeCfg) getCurrentURL() *url.URL {
74✔
565
        cfg.RLock()
74✔
566
        defer cfg.RUnlock()
74✔
567
        return cfg.curURL
74✔
568
}
74✔
569

570
// Returns how long the server should wait before attempting
571
// to solicit a remote leafnode connection.
572
func (cfg *leafNodeCfg) getConnectDelay() time.Duration {
763✔
573
        cfg.RLock()
763✔
574
        delay := cfg.connDelay
763✔
575
        cfg.RUnlock()
763✔
576
        return delay
763✔
577
}
763✔
578

579
// Sets the connect delay.
580
func (cfg *leafNodeCfg) setConnectDelay(delay time.Duration) {
135✔
581
        cfg.Lock()
135✔
582
        cfg.connDelay = delay
135✔
583
        cfg.Unlock()
135✔
584
}
135✔
585

586
// Ensure that non-exported options (used in tests) have
587
// been properly set.
588
func (s *Server) setLeafNodeNonExportedOptions() {
6,574✔
589
        opts := s.getOpts()
6,574✔
590
        s.leafNodeOpts.dialTimeout = opts.LeafNode.dialTimeout
6,574✔
591
        if s.leafNodeOpts.dialTimeout == 0 {
13,147✔
592
                // Use same timeouts as routes for now.
6,573✔
593
                s.leafNodeOpts.dialTimeout = DEFAULT_ROUTE_DIAL
6,573✔
594
        }
6,573✔
595
        s.leafNodeOpts.resolver = opts.LeafNode.resolver
6,574✔
596
        if s.leafNodeOpts.resolver == nil {
13,145✔
597
                s.leafNodeOpts.resolver = net.DefaultResolver
6,571✔
598
        }
6,571✔
599
}
600

601
const sharedSysAccDelay = 250 * time.Millisecond
602

603
// establishHTTPProxyTunnel establishes an HTTP CONNECT tunnel through a proxy server
604
func establishHTTPProxyTunnel(proxyURL, targetHost string, timeout time.Duration, username, password string) (net.Conn, error) {
11✔
605
        proxyAddr, err := url.Parse(proxyURL)
11✔
606
        if err != nil {
11✔
607
                // This should not happen since proxy URL is validated during configuration parsing
×
608
                return nil, fmt.Errorf("unexpected proxy URL parse error (URL was pre-validated): %v", err)
×
609
        }
×
610

611
        // Connect to the proxy server
612
        conn, err := natsDialTimeout("tcp", proxyAddr.Host, timeout)
11✔
613
        if err != nil {
11✔
614
                return nil, fmt.Errorf("failed to connect to proxy: %v", err)
×
615
        }
×
616

617
        // Set deadline for the entire proxy handshake
618
        if err := conn.SetDeadline(time.Now().Add(timeout)); err != nil {
11✔
619
                conn.Close()
×
620
                return nil, fmt.Errorf("failed to set deadline: %v", err)
×
621
        }
×
622

623
        req := &http.Request{
11✔
624
                Method: http.MethodConnect,
11✔
625
                URL:    &url.URL{Opaque: targetHost}, // Opaque is required for CONNECT
11✔
626
                Host:   targetHost,
11✔
627
                Header: make(http.Header),
11✔
628
        }
11✔
629

11✔
630
        // Add proxy authentication if provided
11✔
631
        if username != "" && password != "" {
13✔
632
                req.Header.Set("Proxy-Authorization", "Basic "+base64.StdEncoding.EncodeToString([]byte(username+":"+password)))
2✔
633
        }
2✔
634

635
        if err := req.Write(conn); err != nil {
11✔
636
                conn.Close()
×
637
                return nil, fmt.Errorf("failed to write CONNECT request: %v", err)
×
638
        }
×
639

640
        resp, err := http.ReadResponse(bufio.NewReader(conn), req)
11✔
641
        if err != nil {
11✔
642
                conn.Close()
×
643
                return nil, fmt.Errorf("failed to read proxy response: %v", err)
×
644
        }
×
645

646
        if resp.StatusCode != http.StatusOK {
12✔
647
                resp.Body.Close()
1✔
648
                conn.Close()
1✔
649
                return nil, fmt.Errorf("proxy CONNECT failed: %s", resp.Status)
1✔
650
        }
1✔
651

652
        // Close the response body
653
        resp.Body.Close()
10✔
654

10✔
655
        // Clear the deadline now that we've finished the proxy handshake
10✔
656
        if err := conn.SetDeadline(time.Time{}); err != nil {
10✔
657
                conn.Close()
×
658
                return nil, fmt.Errorf("failed to clear deadline: %v", err)
×
659
        }
×
660

661
        return conn, nil
10✔
662
}
663

664
// Connect to a remote leaf node asynchronously (that is, this function will do
665
// the connect in a go routine).
666
func (s *Server) connectToRemoteLeafNodeAsynchronously(remote *leafNodeCfg, firstConnect bool) {
580✔
667
        remote.setConnectInProgress(true)
580✔
668
        s.startGoRoutine(func() {
1,160✔
669
                defer s.grWG.Done()
580✔
670
                if !connectToRemoteLeafNode(s, remote, firstConnect) {
663✔
671
                        remote.setConnectInProgress(false)
83✔
672
                }
83✔
673
        })
674
}
675

676
// Connect to a remote leaf node. Should only be invoked from
677
// `s.connectToRemoteLeafNodeAsynchronously()` or `s.reConnectToRemoteLeafNode()`.
678
// Returns `true` if this function invoked `s.createLeafNode()`, false otherwise.
679
func connectToRemoteLeafNode(s *Server, remote *leafNodeCfg, firstConnect bool) bool {
763✔
680

763✔
681
        if remote == nil || len(remote.URLs) == 0 {
763✔
682
                s.Debugf("Empty remote leafnode definition, nothing to connect")
×
683
                return false
×
684
        }
×
685

686
        // If this remote is no longer valid by the time we return (e.g. it was
687
        // removed or disabled through a configuration reload), we will never
688
        // reconnect, so clear any JetStream observer state. Otherwise, the raft
689
        // nodes of this account's assets would remain observers.
690
        defer func() {
1,525✔
691
                if remote.stillValid() {
1,522✔
692
                        return
760✔
693
                }
760✔
694
                s.mu.RLock()
2✔
695
                shouldMigrate := remote.JetStreamClusterMigrate
2✔
696
                s.mu.RUnlock()
2✔
697
                if shouldMigrate {
4✔
698
                        s.clearObserverState(remote)
2✔
699
                }
2✔
700
        }()
701

702
        opts := s.getOpts()
763✔
703
        reconnectDelay := opts.LeafNode.ReconnectInterval
763✔
704
        s.mu.RLock()
763✔
705
        dialTimeout := s.leafNodeOpts.dialTimeout
763✔
706
        resolver := s.leafNodeOpts.resolver
763✔
707
        var isSysAcc bool
763✔
708
        if s.eventsEnabled() {
1,493✔
709
                isSysAcc = remote.LocalAccount == s.sys.account.Name
730✔
710
        }
730✔
711
        jetstreamMigrateDelay := remote.JetStreamClusterMigrateDelay
763✔
712
        s.mu.RUnlock()
763✔
713

763✔
714
        // If we are sharing a system account and we are not standalone delay to gather some info prior.
763✔
715
        if firstConnect && isSysAcc && !s.standAloneMode() {
826✔
716
                s.Debugf("Will delay first leafnode connect to shared system account due to clustering")
63✔
717
                remote.setConnectDelay(sharedSysAccDelay)
63✔
718
        }
63✔
719

720
        if connDelay := remote.getConnectDelay(); connDelay > 0 {
833✔
721
                select {
70✔
722
                case <-time.After(connDelay):
62✔
723
                case <-remote.quitCh:
×
724
                        return false
×
725
                case <-s.quitCh:
8✔
726
                        return false
8✔
727
                }
728
                remote.setConnectDelay(0)
62✔
729
        }
730

731
        var conn net.Conn
755✔
732

755✔
733
        const connErrFmt = "Error trying to connect as leafnode to remote server %q (attempt %v): %v"
755✔
734

755✔
735
        // Capture proxy configuration once before the loop with proper locking
755✔
736
        remote.RLock()
755✔
737
        proxyURL := remote.Proxy.URL
755✔
738
        proxyUsername := remote.Proxy.Username
755✔
739
        proxyPassword := remote.Proxy.Password
755✔
740
        proxyTimeout := remote.Proxy.Timeout
755✔
741
        remote.RUnlock()
755✔
742

755✔
743
        // Set default proxy timeout if not specified
755✔
744
        if proxyTimeout == 0 {
1,502✔
745
                proxyTimeout = dialTimeout
747✔
746
        }
747✔
747

748
        attempts := 0
755✔
749

755✔
750
        // In case the migrate timer was created but not canceled, do it when
755✔
751
        // this function exits. Note that the timer would not be created if
755✔
752
        // `jetstreamMigrateDelay == 0`.
755✔
753
        if jetstreamMigrateDelay > 0 {
763✔
754
                defer remote.cancelMigrateTimer()
8✔
755
        }
8✔
756

757
        reconnectTimer := time.NewTimer(reconnectDelay)
755✔
758
        reconnectTimer.Stop()
755✔
759
        defer stopAndClearTimer(&reconnectTimer)
755✔
760

755✔
761
        for s.isRunning() && remote.stillValid() {
3,900✔
762
                rURL := remote.pickNextURL()
3,145✔
763
                url, err := s.getRandomIP(resolver, rURL.Host, nil)
3,145✔
764
                if err == nil {
6,285✔
765
                        var ipStr string
3,140✔
766
                        if url != rURL.Host {
3,201✔
767
                                ipStr = fmt.Sprintf(" (%s)", url)
61✔
768
                        }
61✔
769
                        // Some test may want to disable remotes from connecting
770
                        if s.isLeafConnectDisabled() {
3,276✔
771
                                s.Debugf("Will not attempt to connect to remote server on %q%s, leafnodes currently disabled", rURL.Host, ipStr)
136✔
772
                                err = ErrLeafNodeDisabled
136✔
773
                        } else {
3,140✔
774
                                s.Debugf("Trying to connect as leafnode to remote server on %q%s", rURL.Host, ipStr)
3,004✔
775

3,004✔
776
                                // Check if proxy is configured
3,004✔
777
                                if proxyURL != _EMPTY_ {
3,012✔
778
                                        targetHost := rURL.Host
8✔
779
                                        // If URL doesn't include port, add the default port for the scheme
8✔
780
                                        if rURL.Port() == _EMPTY_ {
8✔
781
                                                defaultPort := "80"
×
782
                                                if rURL.Scheme == wsSchemePrefixTLS {
×
783
                                                        defaultPort = "443"
×
784
                                                }
×
785
                                                targetHost = net.JoinHostPort(rURL.Hostname(), defaultPort)
×
786
                                        }
787

788
                                        conn, err = establishHTTPProxyTunnel(proxyURL, targetHost, proxyTimeout, proxyUsername, proxyPassword)
8✔
789
                                } else {
2,996✔
790
                                        // Direct connection
2,996✔
791
                                        conn, err = natsDialTimeout("tcp", url, dialTimeout)
2,996✔
792
                                }
2,996✔
793
                        }
794
                }
795
                if err != nil {
5,618✔
796
                        jitter := time.Duration(rand.Int63n(int64(reconnectDelay)))
2,473✔
797
                        delay := reconnectDelay + jitter
2,473✔
798
                        attempts++
2,473✔
799
                        if s.shouldReportConnectErr(firstConnect, attempts) {
4,940✔
800
                                s.Errorf(connErrFmt, rURL.Host, attempts, err)
2,467✔
801
                        } else {
2,473✔
802
                                s.Debugf(connErrFmt, rURL.Host, attempts, err)
6✔
803
                        }
6✔
804
                        remote.Lock()
2,473✔
805
                        // if we are using a delay to start migrating assets, kick off a migrate timer.
2,473✔
806
                        if remote.jsMigrateTimer == nil && jetstreamMigrateDelay > 0 {
2,481✔
807
                                remote.jsMigrateTimer = time.AfterFunc(jetstreamMigrateDelay, func() {
16✔
808
                                        s.checkJetStreamMigrate(remote)
8✔
809
                                })
8✔
810
                        }
811
                        remote.Unlock()
2,473✔
812
                        reconnectTimer.Reset(delay)
2,473✔
813
                        select {
2,473✔
814
                        case <-s.quitCh:
75✔
815
                                return false
75✔
816
                        case <-remote.quitCh:
2✔
817
                                return false
2✔
818
                        case <-reconnectTimer.C:
2,395✔
819
                                // Check if we should migrate any JetStream assets immediately while this remote is down.
2,395✔
820
                                // This will be used if JetStreamClusterMigrateDelay was not set
2,395✔
821
                                if jetstreamMigrateDelay == 0 {
4,712✔
822
                                        s.checkJetStreamMigrate(remote)
2,317✔
823
                                }
2,317✔
824
                                continue
2,395✔
825
                        }
826
                }
827
                remote.cancelMigrateTimer()
672✔
828
                // We can check here, but really we will have to check again when the server
672✔
829
                // is about to add to the `s.leafs` map later in the process.
672✔
830
                if !remote.stillValid() {
672✔
831
                        conn.Close()
×
832
                        return false
×
833
                }
×
834

835
                // We have a connection here to a remote server.
836
                // Go ahead and create our leaf node and return.
837
                s.createLeafNode(conn, rURL, remote, nil)
672✔
838

672✔
839
                // Clear any observer states if we had them.
672✔
840
                s.clearObserverState(remote)
672✔
841

672✔
842
                return true
672✔
843
        }
844

845
        return false
5✔
846
}
847

848
func (cfg *leafNodeCfg) cancelMigrateTimer() {
680✔
849
        cfg.Lock()
680✔
850
        stopAndClearTimer(&cfg.jsMigrateTimer)
680✔
851
        cfg.Unlock()
680✔
852
}
680✔
853

854
// This will clear any observer state such that stream or consumer assets on this server can become leaders again.
855
func (s *Server) clearObserverState(remote *leafNodeCfg) {
674✔
856
        s.mu.RLock()
674✔
857
        accName := remote.LocalAccount
674✔
858
        s.mu.RUnlock()
674✔
859

674✔
860
        acc, err := s.LookupAccount(accName)
674✔
861
        if err != nil {
676✔
862
                s.Warnf("Error looking up account [%s] checking for JetStream clear observer state on a leafnode", accName)
2✔
863
                return
2✔
864
        }
2✔
865

866
        acc.jscmMu.Lock()
672✔
867
        defer acc.jscmMu.Unlock()
672✔
868

672✔
869
        // Walk all streams looking for any clustered stream, skip otherwise.
672✔
870
        for _, mset := range acc.streams() {
695✔
871
                node := mset.raftNode()
23✔
872
                if node == nil {
37✔
873
                        // Not R>1
14✔
874
                        continue
14✔
875
                }
876
                // Check consumers
877
                for _, o := range mset.getConsumers() {
11✔
878
                        if n := o.raftNode(); n != nil {
4✔
879
                                // Ensure we can become a leader again.
2✔
880
                                n.SetObserver(false)
2✔
881
                        }
2✔
882
                }
883
                // Ensure we can not become a leader again.
884
                node.SetObserver(false)
9✔
885
        }
886
}
887

888
// Check to see if we should migrate any assets from this account.
889
func (s *Server) checkJetStreamMigrate(remote *leafNodeCfg) {
2,325✔
890
        s.mu.RLock()
2,325✔
891
        accName, shouldMigrate := remote.LocalAccount, remote.JetStreamClusterMigrate
2,325✔
892
        s.mu.RUnlock()
2,325✔
893

2,325✔
894
        if !shouldMigrate {
4,586✔
895
                return
2,261✔
896
        }
2,261✔
897

898
        acc, err := s.LookupAccount(accName)
64✔
899
        if err != nil {
64✔
900
                s.Warnf("Error looking up account [%s] checking for JetStream migration on a leafnode", accName)
×
901
                return
×
902
        }
×
903

904
        acc.jscmMu.Lock()
64✔
905
        defer acc.jscmMu.Unlock()
64✔
906

64✔
907
        // Walk all streams looking for any clustered stream, skip otherwise.
64✔
908
        // If we are the leader force stepdown.
64✔
909
        for _, mset := range acc.streams() {
96✔
910
                node := mset.raftNode()
32✔
911
                if node == nil {
32✔
912
                        // Not R>1
×
913
                        continue
×
914
                }
915
                // Collect any consumers
916
                for _, o := range mset.getConsumers() {
52✔
917
                        if n := o.raftNode(); n != nil {
40✔
918
                                n.StepDown()
20✔
919
                                // Ensure we can not become a leader while in this state.
20✔
920
                                n.SetObserver(true)
20✔
921
                        }
20✔
922
                }
923
                // Stepdown if this stream was leader.
924
                node.StepDown()
32✔
925
                // Ensure we can not become a leader while in this state.
32✔
926
                node.SetObserver(true)
32✔
927
        }
928
}
929

930
// Helper for checking.
931
func (s *Server) isLeafConnectDisabled() bool {
3,140✔
932
        s.mu.RLock()
3,140✔
933
        defer s.mu.RUnlock()
3,140✔
934
        return s.leafDisableConnect
3,140✔
935
}
3,140✔
936

937
// Save off the tlsName for when we use TLS and mix hostnames and IPs. IPs usually
938
// come from the server we connect to.
939
//
940
// We used to save the name only if there was a TLSConfig or scheme equal to "tls".
941
// However, this was causing failures for users that did not set the scheme (and
942
// their remote connections did not have a tls{} block).
943
// We now save the host name regardless in case the remote returns an INFO indicating
944
// that TLS is required.
945
//
946
// Lock held on entry.
947
func (cfg *leafNodeCfg) saveTLSHostname(u *url.URL) {
1,360✔
948
        if cfg.tlsName == _EMPTY_ && net.ParseIP(u.Hostname()) == nil {
1,373✔
949
                cfg.tlsName = u.Hostname()
13✔
950
        }
13✔
951
}
952

953
// Save off the username/password for when we connect using a bare URL
954
// that we get from the INFO protocol.
955
//
956
// Lock held on entry.
957
func (cfg *leafNodeCfg) saveUserPassword(u *url.URL) {
887✔
958
        if cfg.username == _EMPTY_ && u.User != nil {
1,109✔
959
                cfg.username = u.User.Username()
222✔
960
                cfg.password, _ = u.User.Password()
222✔
961
        }
222✔
962
}
963

964
// This starts the leafnode accept loop in a go routine, unless it
965
// is detected that the server has already been shutdown.
966
func (s *Server) startLeafNodeAcceptLoop() {
3,982✔
967
        // Snapshot server options.
3,982✔
968
        opts := s.getOpts()
3,982✔
969

3,982✔
970
        port := opts.LeafNode.Port
3,982✔
971
        if port == -1 {
7,815✔
972
                port = 0
3,833✔
973
        }
3,833✔
974

975
        if s.isShuttingDown() {
3,982✔
976
                return
×
977
        }
×
978

979
        s.mu.Lock()
3,982✔
980
        hp := net.JoinHostPort(opts.LeafNode.Host, strconv.Itoa(port))
3,982✔
981
        l, e := natsListen("tcp", hp)
3,982✔
982
        s.leafNodeListenerErr = e
3,982✔
983
        if e != nil {
3,982✔
984
                s.mu.Unlock()
×
985
                s.Fatalf("Error listening on leafnode port: %d - %v", opts.LeafNode.Port, e)
×
986
                return
×
987
        }
×
988

989
        s.Noticef("Listening for leafnode connections on %s",
3,982✔
990
                net.JoinHostPort(opts.LeafNode.Host, strconv.Itoa(l.Addr().(*net.TCPAddr).Port)))
3,982✔
991

3,982✔
992
        tlsRequired := opts.LeafNode.TLSConfig != nil
3,982✔
993
        tlsVerify := tlsRequired && opts.LeafNode.TLSConfig.ClientAuth == tls.RequireAndVerifyClientCert
3,982✔
994
        // Do not set compression in this Info object, it would possibly cause
3,982✔
995
        // issues when sending asynchronous INFO to the remote.
3,982✔
996
        info := Info{
3,982✔
997
                ID:            s.info.ID,
3,982✔
998
                Name:          s.info.Name,
3,982✔
999
                Version:       s.info.Version,
3,982✔
1000
                GitCommit:     gitCommit,
3,982✔
1001
                GoVersion:     runtime.Version(),
3,982✔
1002
                AuthRequired:  true,
3,982✔
1003
                TLSRequired:   tlsRequired,
3,982✔
1004
                TLSVerify:     tlsVerify,
3,982✔
1005
                MaxPayload:    s.info.MaxPayload, // TODO(dlc) - Allow override?
3,982✔
1006
                Headers:       s.supportsHeaders(),
3,982✔
1007
                JetStream:     opts.JetStream,
3,982✔
1008
                Domain:        opts.JetStreamDomain,
3,982✔
1009
                Proto:         s.getServerProto(),
3,982✔
1010
                InfoOnConnect: true,
3,982✔
1011
                JSApiLevel:    JSApiLevel,
3,982✔
1012
        }
3,982✔
1013
        // If we have selected a random port...
3,982✔
1014
        if port == 0 {
7,815✔
1015
                // Write resolved port back to options.
3,833✔
1016
                opts.LeafNode.Port = l.Addr().(*net.TCPAddr).Port
3,833✔
1017
        }
3,833✔
1018

1019
        s.leafNodeInfo = info
3,982✔
1020
        // Possibly override Host/Port and set IP based on Cluster.Advertise
3,982✔
1021
        if err := s.setLeafNodeInfoHostPortAndIP(); err != nil {
3,982✔
1022
                s.Fatalf("Error setting leafnode INFO with LeafNode.Advertise value of %s, err=%v", opts.LeafNode.Advertise, err)
×
1023
                l.Close()
×
1024
                s.mu.Unlock()
×
1025
                return
×
1026
        }
×
1027
        s.leafURLsMap[s.leafNodeInfo.IP]++
3,982✔
1028
        s.generateLeafNodeInfoJSON()
3,982✔
1029

3,982✔
1030
        // Setup state that can enable shutdown
3,982✔
1031
        s.leafNodeListener = l
3,982✔
1032

3,982✔
1033
        // As of now, a server that does not have remotes configured would
3,982✔
1034
        // never solicit a connection, so we should not have to warn if
3,982✔
1035
        // InsecureSkipVerify is set in main LeafNodes config (since
3,982✔
1036
        // this TLS setting matters only when soliciting a connection).
3,982✔
1037
        // Still, warn if insecure is set in any of LeafNode block.
3,982✔
1038
        // We need to check remotes, even if tls is not required on accept.
3,982✔
1039
        warn := tlsRequired && opts.LeafNode.TLSConfig.InsecureSkipVerify
3,982✔
1040
        if !warn {
7,962✔
1041
                for _, r := range opts.LeafNode.Remotes {
4,104✔
1042
                        if r.TLSConfig != nil && r.TLSConfig.InsecureSkipVerify {
124✔
1043
                                warn = true
×
1044
                                break
×
1045
                        }
1046
                }
1047
        }
1048
        if warn {
3,984✔
1049
                s.Warnf(leafnodeTLSInsecureWarning)
2✔
1050
        }
2✔
1051
        go s.acceptConnections(l, "Leafnode", func(conn net.Conn) { s.createLeafNode(conn, nil, nil, nil) }, nil)
4,689✔
1052
        s.mu.Unlock()
3,982✔
1053
}
1054

1055
// RegEx to match a creds file with user JWT and Seed.
1056
var credsRe = regexp.MustCompile(`\s*(?:(?:[-]{3,}.*[-]{3,}\r?\n)([\w\-.=]+)(?:\r?\n[-]{3,}.*[-]{3,}(\r?\n|\z)))`)
1057

1058
// clusterName is provided as argument to avoid lock ordering issues with the locked client c
1059
// Lock should be held entering here.
1060
func (c *client) sendLeafConnect(clusterName string, headers bool) error {
548✔
1061
        // We support basic user/pass and operator based user JWT with signatures.
548✔
1062
        cinfo := leafConnectInfo{
548✔
1063
                Version:       VERSION,
548✔
1064
                ID:            c.srv.info.ID,
548✔
1065
                Domain:        c.srv.info.Domain,
548✔
1066
                Name:          c.srv.info.Name,
548✔
1067
                Hub:           c.leaf.remote.Hub,
548✔
1068
                Cluster:       clusterName,
548✔
1069
                Headers:       headers,
548✔
1070
                JetStream:     c.acc.jetStreamConfigured(),
548✔
1071
                DenyPub:       c.leaf.remote.DenyImports,
548✔
1072
                Compression:   c.leaf.compression,
548✔
1073
                RemoteAccount: c.acc.GetName(),
548✔
1074
                Proto:         c.srv.getServerProto(),
548✔
1075
                Isolate:       c.leaf.remote.RequestIsolation,
548✔
1076
        }
548✔
1077

548✔
1078
        // If a signature callback is specified, this takes precedence over anything else.
548✔
1079
        if cb := c.leaf.remote.SignatureCB; cb != nil {
553✔
1080
                nonce := c.nonce
5✔
1081
                c.mu.Unlock()
5✔
1082
                jwt, sigraw, err := cb(nonce)
5✔
1083
                c.mu.Lock()
5✔
1084
                if err == nil && c.isClosed() {
6✔
1085
                        err = ErrConnectionClosed
1✔
1086
                }
1✔
1087
                if err != nil {
7✔
1088
                        c.Errorf("Error signing the nonce: %v", err)
2✔
1089
                        return err
2✔
1090
                }
2✔
1091
                sig := base64.RawURLEncoding.EncodeToString(sigraw)
3✔
1092
                cinfo.JWT, cinfo.Sig = jwt, sig
3✔
1093

1094
        } else if creds := c.leaf.remote.Credentials; creds != _EMPTY_ {
587✔
1095
                // Check for credentials first, that will take precedence..
44✔
1096
                c.Debugf("Authenticating with credentials file %q", c.leaf.remote.Credentials)
44✔
1097
                contents, err := os.ReadFile(creds)
44✔
1098
                if err != nil {
44✔
1099
                        c.Errorf("%v", err)
×
1100
                        return err
×
1101
                }
×
1102
                defer wipeSlice(contents)
44✔
1103
                items := credsRe.FindAllSubmatch(contents, -1)
44✔
1104
                if len(items) < 2 {
44✔
1105
                        c.Errorf("Credentials file malformed")
×
1106
                        return err
×
1107
                }
×
1108
                // First result should be the user JWT.
1109
                // We copy here so that the file containing the seed will be wiped appropriately.
1110
                raw := items[0][1]
44✔
1111
                tmp := make([]byte, len(raw))
44✔
1112
                copy(tmp, raw)
44✔
1113
                // Seed is second item.
44✔
1114
                kp, err := nkeys.FromSeed(items[1][1])
44✔
1115
                if err != nil {
44✔
1116
                        c.Errorf("Credentials file has malformed seed")
×
1117
                        return err
×
1118
                }
×
1119
                // Wipe our key on exit.
1120
                defer kp.Wipe()
44✔
1121

44✔
1122
                sigraw, _ := kp.Sign(c.nonce)
44✔
1123
                sig := base64.RawURLEncoding.EncodeToString(sigraw)
44✔
1124
                cinfo.JWT = bytesToString(tmp)
44✔
1125
                cinfo.Sig = sig
44✔
1126
        } else if nkey := c.leaf.remote.Nkey; nkey != _EMPTY_ {
502✔
1127
                kp, err := nkeys.FromSeed([]byte(nkey))
3✔
1128
                if err != nil {
3✔
1129
                        c.Errorf("Remote nkey has malformed seed")
×
1130
                        return err
×
1131
                }
×
1132
                // Wipe our key on exit.
1133
                defer kp.Wipe()
3✔
1134
                sigraw, _ := kp.Sign(c.nonce)
3✔
1135
                sig := base64.RawURLEncoding.EncodeToString(sigraw)
3✔
1136
                pkey, _ := kp.PublicKey()
3✔
1137
                cinfo.Nkey = pkey
3✔
1138
                cinfo.Sig = sig
3✔
1139
        }
1140
        // In addition, and this is to allow auth callout, set user/password or
1141
        // token if applicable.
1142
        if userInfo := c.leaf.remote.curURL.User; userInfo != nil {
789✔
1143
                cinfo.User = userInfo.Username()
243✔
1144
                var ok bool
243✔
1145
                cinfo.Pass, ok = userInfo.Password()
243✔
1146
                // For backward compatibility, if only username is provided, set both
243✔
1147
                // Token and User, not just Token.
243✔
1148
                if !ok {
251✔
1149
                        cinfo.Token = cinfo.User
8✔
1150
                }
8✔
1151
        } else if c.leaf.remote.username != _EMPTY_ {
309✔
1152
                cinfo.User = c.leaf.remote.username
6✔
1153
                cinfo.Pass = c.leaf.remote.password
6✔
1154
                // For backward compatibility, if only username is provided, set both
6✔
1155
                // Token and User, not just Token.
6✔
1156
                if cinfo.Pass == _EMPTY_ {
6✔
1157
                        cinfo.Token = cinfo.User
×
1158
                }
×
1159
        }
1160
        b, err := json.Marshal(cinfo)
546✔
1161
        if err != nil {
546✔
1162
                c.Errorf("Error marshaling CONNECT to remote leafnode: %v\n", err)
×
1163
                return err
×
1164
        }
×
1165
        // Although this call is made before the writeLoop is created,
1166
        // we don't really need to send in place. The protocol will be
1167
        // sent out by the writeLoop.
1168
        c.enqueueProto([]byte(fmt.Sprintf(ConProto, b)))
546✔
1169
        return nil
546✔
1170
}
1171

1172
// Makes a deep copy of the LeafNode Info structure.
1173
// The server lock is held on entry.
1174
func (s *Server) copyLeafNodeInfo() *Info {
2,183✔
1175
        clone := s.leafNodeInfo
2,183✔
1176
        // Copy the array of urls.
2,183✔
1177
        if len(s.leafNodeInfo.LeafNodeURLs) > 0 {
3,965✔
1178
                clone.LeafNodeURLs = append([]string(nil), s.leafNodeInfo.LeafNodeURLs...)
1,782✔
1179
        }
1,782✔
1180
        return &clone
2,183✔
1181
}
1182

1183
// Adds a LeafNode URL that we get when a route connects to the Info structure.
1184
// Regenerates the JSON byte array so that it can be sent to LeafNode connections.
1185
// Returns a boolean indicating if the URL was added or not.
1186
// Server lock is held on entry
1187
func (s *Server) addLeafNodeURL(urlStr string) bool {
8,243✔
1188
        if s.leafURLsMap.addUrl(urlStr) {
16,481✔
1189
                s.generateLeafNodeInfoJSON()
8,238✔
1190
                return true
8,238✔
1191
        }
8,238✔
1192
        return false
5✔
1193
}
1194

1195
// Removes a LeafNode URL of the route that is disconnecting from the Info structure.
1196
// Regenerates the JSON byte array so that it can be sent to LeafNode connections.
1197
// Returns a boolean indicating if the URL was removed or not.
1198
// Server lock is held on entry.
1199
func (s *Server) removeLeafNodeURL(urlStr string) bool {
8,243✔
1200
        // Don't need to do this if we are removing the route connection because
8,243✔
1201
        // we are shuting down...
8,243✔
1202
        if s.isShuttingDown() {
12,633✔
1203
                return false
4,390✔
1204
        }
4,390✔
1205
        if s.leafURLsMap.removeUrl(urlStr) {
7,702✔
1206
                s.generateLeafNodeInfoJSON()
3,849✔
1207
                return true
3,849✔
1208
        }
3,849✔
1209
        return false
4✔
1210
}
1211

1212
// Server lock is held on entry
1213
func (s *Server) generateLeafNodeInfoJSON() {
16,069✔
1214
        s.leafNodeInfo.Cluster = s.cachedClusterName()
16,069✔
1215
        s.leafNodeInfo.LeafNodeURLs = s.leafURLsMap.getAsStringSlice()
16,069✔
1216
        s.leafNodeInfo.WSConnectURLs = s.websocket.connectURLsMap.getAsStringSlice()
16,069✔
1217
        s.leafNodeInfoJSON = generateInfoJSON(&s.leafNodeInfo)
16,069✔
1218
}
16,069✔
1219

1220
// Sends an async INFO protocol so that the connected servers can update
1221
// their list of LeafNode urls.
1222
func (s *Server) sendAsyncLeafNodeInfo() {
12,087✔
1223
        for _, c := range s.leafs {
12,146✔
1224
                c.mu.Lock()
59✔
1225
                c.enqueueProto(s.leafNodeInfoJSON)
59✔
1226
                c.mu.Unlock()
59✔
1227
        }
59✔
1228
}
1229

1230
// Called when an inbound leafnode connection is accepted or we create one for a solicited leafnode.
1231
func (s *Server) createLeafNode(conn net.Conn, rURL *url.URL, remote *leafNodeCfg, ws *websocket) *client {
1,405✔
1232
        // Snapshot server options.
1,405✔
1233
        opts := s.getOpts()
1,405✔
1234

1,405✔
1235
        maxPay := int32(opts.MaxPayload)
1,405✔
1236
        maxSubs := int32(opts.MaxSubs)
1,405✔
1237
        // For system, maxSubs of 0 means unlimited, so re-adjust here.
1,405✔
1238
        if maxSubs == 0 {
2,809✔
1239
                maxSubs = -1
1,404✔
1240
        }
1,404✔
1241
        now := time.Now().UTC()
1,405✔
1242

1,405✔
1243
        c := &client{srv: s, nc: conn, kind: LEAF, opts: defaultOpts, mpay: maxPay, msubs: maxSubs, start: now, last: now}
1,405✔
1244
        // Do not update the smap here, we need to do it in initLeafNodeSmapAndSendSubs
1,405✔
1245
        c.leaf = &leaf{}
1,405✔
1246

1,405✔
1247
        // If the leafnode subject interest should be isolated, flag it here.
1,405✔
1248
        s.optsMu.RLock()
1,405✔
1249
        if c.leaf.isolated = s.opts.LeafNode.IsolateLeafnodeInterest; !c.leaf.isolated && remote != nil {
2,077✔
1250
                c.leaf.isolated = remote.LocalIsolation
672✔
1251
        }
672✔
1252
        s.optsMu.RUnlock()
1,405✔
1253

1,405✔
1254
        // For accepted LN connections, ws will be != nil if it was accepted
1,405✔
1255
        // through the Websocket port.
1,405✔
1256
        c.ws = ws
1,405✔
1257

1,405✔
1258
        // For remote, check if the scheme starts with "ws", if so, we will initiate
1,405✔
1259
        // a remote Leaf Node connection as a websocket connection.
1,405✔
1260
        if remote != nil && rURL != nil && isWSURL(rURL) {
1,451✔
1261
                remote.RLock()
46✔
1262
                c.ws = &websocket{compress: remote.Websocket.Compression, maskwrite: !remote.Websocket.NoMasking}
46✔
1263
                remote.RUnlock()
46✔
1264
        }
46✔
1265

1266
        // Determines if we are soliciting the connection or not.
1267
        var solicited bool
1,405✔
1268
        var acc *Account
1,405✔
1269
        var remoteSuffix string
1,405✔
1270
        if remote != nil {
2,077✔
1271
                // For now, if lookup fails, we will constantly try
672✔
1272
                // to recreate this LN connection.
672✔
1273
                lacc := remote.LocalAccount
672✔
1274
                var err error
672✔
1275
                acc, err = s.LookupAccount(lacc)
672✔
1276
                if err != nil {
674✔
1277
                        // An account not existing is something that can happen with nats/http account resolver and the account
2✔
1278
                        // has not yet been pushed, or the request failed for other reasons.
2✔
1279
                        // remote needs to be set or retry won't happen
2✔
1280
                        c.leaf.remote = remote
2✔
1281
                        c.closeConnection(MissingAccount)
2✔
1282
                        s.Errorf("Unable to lookup account %s for solicited leafnode connection: %v", lacc, err)
2✔
1283
                        return nil
2✔
1284
                }
2✔
1285
                remoteSuffix = fmt.Sprintf(" for account: %s", acc.traceLabel())
670✔
1286
        }
1287

1288
        c.mu.Lock()
1,403✔
1289
        c.initClient()
1,403✔
1290
        c.Noticef("Leafnode connection created%s %s", remoteSuffix, c.opts.Name)
1,403✔
1291

1,403✔
1292
        var (
1,403✔
1293
                tlsFirst         bool
1,403✔
1294
                tlsFirstFallback time.Duration
1,403✔
1295
                infoTimeout      time.Duration
1,403✔
1296
        )
1,403✔
1297
        if remote != nil {
2,073✔
1298
                solicited = true
670✔
1299
                remote.Lock()
670✔
1300
                c.leaf.remote = remote
670✔
1301
                c.setPermissions(remote.perms)
670✔
1302
                if !c.leaf.remote.Hub {
1,328✔
1303
                        c.leaf.isSpoke = true
658✔
1304
                }
658✔
1305
                tlsFirst = remote.TLSHandshakeFirst
670✔
1306
                infoTimeout = remote.FirstInfoTimeout
670✔
1307
                remote.Unlock()
670✔
1308
                c.acc = acc
670✔
1309
        } else {
733✔
1310
                c.flags.set(expectConnect)
733✔
1311
                if ws != nil {
759✔
1312
                        c.Debugf("Leafnode compression=%v", c.ws.compress)
26✔
1313
                }
26✔
1314
                tlsFirst = opts.LeafNode.TLSHandshakeFirst
733✔
1315
                if f := opts.LeafNode.TLSHandshakeFirstFallback; f > 0 {
734✔
1316
                        tlsFirstFallback = f
1✔
1317
                }
1✔
1318
        }
1319
        c.mu.Unlock()
1,403✔
1320

1,403✔
1321
        var nonce [nonceLen]byte
1,403✔
1322
        var info *Info
1,403✔
1323

1,403✔
1324
        // Grab this before the client lock below.
1,403✔
1325
        if !solicited {
2,136✔
1326
                // Grab server variables
733✔
1327
                s.mu.Lock()
733✔
1328
                info = s.copyLeafNodeInfo()
733✔
1329
                // For tests that want to simulate old servers, do not set the compression
733✔
1330
                // on the INFO protocol if configured with CompressionNotSupported.
733✔
1331
                // Also suppress it if WebSocket compression is already in use, otherwise
733✔
1332
                // an old soliciting peer would honor the advertised mode, switch to S2,
733✔
1333
                // and then wait forever for a compressed INFO response from us.
733✔
1334
                if cm := opts.LeafNode.Compression.Mode; cm != CompressionNotSupported && (ws == nil || !ws.compress) {
1,462✔
1335
                        info.Compression = cm
729✔
1336
                }
729✔
1337
                // We always send a nonce for LEAF connections. Do not change that without
1338
                // taking into account presence of proxy trusted keys.
1339
                s.generateNonce(nonce[:])
733✔
1340
                s.mu.Unlock()
733✔
1341
        }
1342

1343
        // Grab lock
1344
        c.mu.Lock()
1,403✔
1345

1,403✔
1346
        var preBuf []byte
1,403✔
1347
        if solicited {
2,073✔
1348
                // For websocket connection, we need to send an HTTP request,
670✔
1349
                // and get the response before starting the readLoop to get
670✔
1350
                // the INFO, etc..
670✔
1351
                if c.isWebsocket() {
716✔
1352
                        var err error
46✔
1353
                        var closeReason ClosedState
46✔
1354

46✔
1355
                        preBuf, closeReason, err = c.leafNodeSolicitWSConnection(opts, rURL, remote)
46✔
1356
                        if err != nil {
66✔
1357
                                c.Errorf("Error soliciting websocket connection: %v", err)
20✔
1358
                                c.mu.Unlock()
20✔
1359
                                if closeReason != 0 {
36✔
1360
                                        c.closeConnection(closeReason)
16✔
1361
                                }
16✔
1362
                                return nil
20✔
1363
                        }
1364
                } else {
624✔
1365
                        // If configured to do TLS handshake first
624✔
1366
                        if tlsFirst {
628✔
1367
                                if _, err := c.leafClientHandshakeIfNeeded(remote, opts); err != nil {
5✔
1368
                                        c.mu.Unlock()
1✔
1369
                                        return nil
1✔
1370
                                }
1✔
1371
                        }
1372
                        // We need to wait for the info, but not for too long.
1373
                        c.nc.SetReadDeadline(time.Now().Add(infoTimeout))
623✔
1374
                }
1375

1376
                // We will process the INFO from the readloop and finish by
1377
                // sending the CONNECT and finish registration later.
1378
        } else {
733✔
1379
                // Send our info to the other side.
733✔
1380
                // Remember the nonce we sent here for signatures, etc.
733✔
1381
                c.nonce = make([]byte, nonceLen)
733✔
1382
                copy(c.nonce, nonce[:])
733✔
1383
                info.Nonce = bytesToString(c.nonce)
733✔
1384
                info.CID = c.cid
733✔
1385
                proto := generateInfoJSON(info)
733✔
1386

733✔
1387
                var pre []byte
733✔
1388
                // We need first to check for "TLS First" fallback delay.
733✔
1389
                if tlsFirstFallback > 0 {
734✔
1390
                        // We wait and see if we are getting any data. Since we did not send
1✔
1391
                        // the INFO protocol yet, only clients that use TLS first should be
1✔
1392
                        // sending data (the TLS handshake). We don't really check the content:
1✔
1393
                        // if it is a rogue agent and not an actual client performing the
1✔
1394
                        // TLS handshake, the error will be detected when performing the
1✔
1395
                        // handshake on our side.
1✔
1396
                        pre = make([]byte, 4)
1✔
1397
                        c.nc.SetReadDeadline(time.Now().Add(tlsFirstFallback))
1✔
1398
                        n, _ := io.ReadFull(c.nc, pre[:])
1✔
1399
                        c.nc.SetReadDeadline(time.Time{})
1✔
1400
                        // If we get any data (regardless of possible timeout), we will proceed
1✔
1401
                        // with the TLS handshake.
1✔
1402
                        if n > 0 {
1✔
1403
                                pre = pre[:n]
×
1404
                        } else {
1✔
1405
                                // We did not get anything so we will send the INFO protocol.
1✔
1406
                                pre = nil
1✔
1407
                                // Set the boolean to false for the rest of the function.
1✔
1408
                                tlsFirst = false
1✔
1409
                        }
1✔
1410
                }
1411

1412
                if !tlsFirst {
1,461✔
1413
                        // We have to send from this go routine because we may
728✔
1414
                        // have to block for TLS handshake before we start our
728✔
1415
                        // writeLoop go routine. The other side needs to receive
728✔
1416
                        // this before it can initiate the TLS handshake..
728✔
1417
                        c.sendProtoNow(proto)
728✔
1418

728✔
1419
                        // The above call could have marked the connection as closed (due to TCP error).
728✔
1420
                        if c.isClosed() {
728✔
1421
                                c.mu.Unlock()
×
1422
                                c.closeConnection(WriteError)
×
1423
                                return nil
×
1424
                        }
×
1425
                }
1426

1427
                // Check to see if we need to spin up TLS.
1428
                if !c.isWebsocket() && info.TLSRequired {
805✔
1429
                        // If we have a prebuffer create a multi-reader.
72✔
1430
                        if len(pre) > 0 {
72✔
1431
                                c.nc = &tlsMixConn{c.nc, bytes.NewBuffer(pre)}
×
1432
                        }
×
1433
                        // Perform server-side TLS handshake.
1434
                        if err := c.doTLSServerHandshake(tlsHandshakeLeaf, opts.LeafNode.TLSConfig, opts.LeafNode.TLSTimeout, opts.LeafNode.TLSPinnedCerts); err != nil {
117✔
1435
                                c.mu.Unlock()
45✔
1436
                                return nil
45✔
1437
                        }
45✔
1438
                }
1439

1440
                // If the user wants the TLS handshake to occur first, now that it is
1441
                // done, send the INFO protocol.
1442
                if tlsFirst {
691✔
1443
                        c.flags.set(didTLSFirst)
3✔
1444
                        c.sendProtoNow(proto)
3✔
1445
                        if c.isClosed() {
3✔
1446
                                c.mu.Unlock()
×
1447
                                c.closeConnection(WriteError)
×
1448
                                return nil
×
1449
                        }
×
1450
                }
1451

1452
                // Leaf nodes will always require a CONNECT to let us know
1453
                // when we are properly bound to an account.
1454
                //
1455
                // If compression is configured, we can't set the authTimer here because
1456
                // it would cause the parser to fail any incoming protocol that is not a
1457
                // CONNECT (and we need to exchange INFO protocols for compression
1458
                // negotiation). So instead, use the ping timer until we are done with
1459
                // negotiation and can set the auth timer.
1460
                timeout := secondsToDuration(opts.LeafNode.AuthTimeout)
688✔
1461
                if needsCompression(opts.LeafNode.Compression.Mode) {
1,167✔
1462
                        c.ping.tmr = time.AfterFunc(timeout, func() {
479✔
1463
                                c.authTimeout()
×
1464
                        })
×
1465
                } else {
209✔
1466
                        c.setAuthTimer(timeout)
209✔
1467
                }
209✔
1468
        }
1469

1470
        // Keep track in case server is shutdown before we can successfully register.
1471
        if !s.addToTempClients(c.cid, c) {
1,338✔
1472
                c.mu.Unlock()
1✔
1473
                c.setNoReconnect()
1✔
1474
                c.closeConnection(ServerShutdown)
1✔
1475
                return nil
1✔
1476
        }
1✔
1477

1478
        // Spin up the read loop.
1479
        s.startGoRoutine(func() { c.readLoop(preBuf) })
2,672✔
1480

1481
        // We will spin the write loop for solicited connections only
1482
        // when processing the INFO and after switching to TLS if needed.
1483
        if !solicited {
2,024✔
1484
                s.startGoRoutine(func() { c.writeLoop() })
1,376✔
1485
        }
1486

1487
        c.mu.Unlock()
1,336✔
1488

1,336✔
1489
        return c
1,336✔
1490
}
1491

1492
// Will perform the client-side TLS handshake if needed. Assumes that this
1493
// is called by the solicit side (remote will be non nil). Returns `true`
1494
// if TLS is required, `false` otherwise.
1495
// Lock held on entry.
1496
func (c *client) leafClientHandshakeIfNeeded(remote *leafNodeCfg, opts *Options) (bool, error) {
1,536✔
1497
        // Check if TLS is required and gather TLS config variables.
1,536✔
1498
        tlsRequired, tlsConfig, tlsName, tlsTimeout := c.leafNodeGetTLSConfigForSolicit(remote)
1,536✔
1499
        if !tlsRequired {
2,998✔
1500
                return false, nil
1,462✔
1501
        }
1,462✔
1502

1503
        // If TLS required, peform handshake.
1504
        // Get the URL that was used to connect to the remote server.
1505
        rURL := remote.getCurrentURL()
74✔
1506

74✔
1507
        // Perform the client-side TLS handshake.
74✔
1508
        if resetTLSName, err := c.doTLSClientHandshake(tlsHandshakeLeaf, rURL, tlsConfig, tlsName, tlsTimeout, opts.LeafNode.TLSPinnedCerts); err != nil {
107✔
1509
                // Check if we need to reset the remote's TLS name.
33✔
1510
                if resetTLSName {
33✔
1511
                        remote.Lock()
×
1512
                        remote.tlsName = _EMPTY_
×
1513
                        remote.Unlock()
×
1514
                }
×
1515
                return false, err
33✔
1516
        }
1517
        return true, nil
41✔
1518
}
1519

1520
func (c *client) processLeafnodeInfo(info *Info) {
2,101✔
1521
        c.mu.Lock()
2,101✔
1522
        if c.leaf == nil || c.isClosed() {
2,101✔
1523
                c.mu.Unlock()
×
1524
                return
×
1525
        }
×
1526
        s := c.srv
2,101✔
1527
        opts := s.getOpts()
2,101✔
1528
        remote := c.leaf.remote
2,101✔
1529
        didSolicit := remote != nil
2,101✔
1530
        firstINFO := !c.flags.isSet(infoReceived)
2,101✔
1531

2,101✔
1532
        // In case of websocket, the TLS handshake has been already done.
2,101✔
1533
        // So check only for non websocket connections and for configurations
2,101✔
1534
        // where the TLS Handshake was not done first.
2,101✔
1535
        if didSolicit && !c.flags.isSet(handshakeComplete) && !c.isWebsocket() && !remote.TLSHandshakeFirst {
3,587✔
1536
                // If the server requires TLS, we need to set this in the remote
1,486✔
1537
                // otherwise if there is no TLS configuration block for the remote,
1,486✔
1538
                // the solicit side will not attempt to perform the TLS handshake.
1,486✔
1539
                if firstINFO && info.TLSRequired {
1,544✔
1540
                        // Check for TLS/proxy configuration mismatch
58✔
1541
                        if remote.Proxy.URL != _EMPTY_ && !remote.TLS && remote.TLSConfig == nil {
58✔
1542
                                c.mu.Unlock()
×
1543
                                c.Errorf("TLS configuration mismatch: Hub requires TLS but leafnode remote is not configured for TLS. When using a proxy, ensure TLS leafnode configuration matches the Hub requirements.")
×
1544
                                c.closeConnection(TLSHandshakeError)
×
1545
                                return
×
1546
                        }
×
1547
                        remote.TLS = true
58✔
1548
                }
1549
                if _, err := c.leafClientHandshakeIfNeeded(remote, opts); err != nil {
1,514✔
1550
                        c.mu.Unlock()
28✔
1551
                        return
28✔
1552
                }
28✔
1553
        }
1554

1555
        // Check for compression, unless already done.
1556
        if firstINFO && !c.flags.isSet(compressionNegotiated) {
3,128✔
1557
                // A solicited leafnode connection must first receive a leafnode INFO.
1,055✔
1558
                // Classify wrong-port connections before any leaf-specific negotiation.
1,055✔
1559
                if didSolicit && (info.CID == 0 || info.LeafNodeURLs == nil) {
1,109✔
1560
                        c.mu.Unlock()
54✔
1561
                        c.Errorf(ErrConnectedToWrongPort.Error())
54✔
1562
                        c.closeConnection(WrongPort)
54✔
1563
                        return
54✔
1564
                }
54✔
1565

1566
                // Prevent from getting back here.
1567
                c.flags.set(compressionNegotiated)
1,001✔
1568

1,001✔
1569
                var co *CompressionOpts
1,001✔
1570
                if !didSolicit {
1,450✔
1571
                        co = &opts.LeafNode.Compression
449✔
1572
                } else {
1,001✔
1573
                        co = &remote.Compression
552✔
1574
                }
552✔
1575
                if needsCompression(co.Mode) {
1,993✔
1576
                        // Release client lock since following function will need server lock.
992✔
1577
                        c.mu.Unlock()
992✔
1578
                        compress, err := s.negotiateLeafCompression(c, didSolicit, info.Compression, co)
992✔
1579
                        if err != nil {
992✔
1580
                                c.sendErrAndErr(err.Error())
×
1581
                                c.closeConnection(ProtocolViolation)
×
1582
                                return
×
1583
                        }
×
1584
                        if compress {
1,892✔
1585
                                // Done for now, will get back another INFO protocol...
900✔
1586
                                return
900✔
1587
                        }
900✔
1588
                        // No compression because one side does not want/can't, so proceed.
1589
                        c.mu.Lock()
92✔
1590
                        // Check that the connection did not close if the lock was released.
92✔
1591
                        if c.isClosed() {
92✔
1592
                                c.mu.Unlock()
×
1593
                                return
×
1594
                        }
×
1595
                } else {
9✔
1596
                        // Coming from an old server, the Compression field would be the empty
9✔
1597
                        // string. For servers that are configured with CompressionNotSupported,
9✔
1598
                        // this makes them behave as old servers.
9✔
1599
                        if info.Compression == _EMPTY_ || co.Mode == CompressionNotSupported {
11✔
1600
                                c.leaf.compression = CompressionNotSupported
2✔
1601
                        } else {
9✔
1602
                                c.leaf.compression = CompressionOff
7✔
1603
                        }
7✔
1604
                }
1605
                // Accepting side does not normally process an INFO protocol during
1606
                // initial connection handshake. So we keep it consistent by returning
1607
                // if we are not soliciting.
1608
                if !didSolicit {
102✔
1609
                        // If we had created the ping timer instead of the auth timer, we will
1✔
1610
                        // clear the ping timer and set the auth timer now that the compression
1✔
1611
                        // negotiation is done.
1✔
1612
                        if info.Compression != _EMPTY_ && c.ping.tmr != nil {
1✔
1613
                                clearTimer(&c.ping.tmr)
×
1614
                                c.setAuthTimer(secondsToDuration(opts.LeafNode.AuthTimeout))
×
1615
                        }
×
1616
                        c.mu.Unlock()
1✔
1617
                        return
1✔
1618
                }
1619
                // Fall through and process the INFO protocol as usual.
1620
        }
1621

1622
        // Note: For now, only the initial INFO has a nonce. We
1623
        // will probably do auto key rotation at some point.
1624
        if firstINFO {
1,688✔
1625
                // Mark that the INFO protocol has been received.
570✔
1626
                c.flags.set(infoReceived)
570✔
1627
                // Prevent connecting to non leafnode port. Need to do this only for
570✔
1628
                // the first INFO, not for async INFO updates...
570✔
1629
                //
570✔
1630
                // Content of INFO sent by the server when accepting a tcp connection.
570✔
1631
                // -------------------------------------------------------------------
570✔
1632
                // Listen Port Of | CID | ClientConnectURLs | LeafNodeURLs | Gateway |
570✔
1633
                // -------------------------------------------------------------------
570✔
1634
                //      CLIENT    |  X* |        X**        |              |         |
570✔
1635
                //      ROUTE     |     |        X**        |      X***    |         |
570✔
1636
                //     GATEWAY    |     |                   |              |    X    |
570✔
1637
                //     LEAFNODE   |  X  |                   |       X      |         |
570✔
1638
                // -------------------------------------------------------------------
570✔
1639
                // *   Not on older servers.
570✔
1640
                // **  Not if "no advertise" is enabled.
570✔
1641
                // *** Not if leafnode's "no advertise" is enabled.
570✔
1642
                //
570✔
1643
                // Reject a cluster that contains spaces.
570✔
1644
                if info.Cluster != _EMPTY_ && strings.Contains(info.Cluster, " ") {
571✔
1645
                        c.mu.Unlock()
1✔
1646
                        c.sendErrAndErr(ErrClusterNameHasSpaces.Error())
1✔
1647
                        c.closeConnection(ProtocolViolation)
1✔
1648
                        return
1✔
1649
                }
1✔
1650
                // For solicited outbound leaf connections, capture the remote's nonce.
1651
                // For inbound leaf connections, keep using the server-issued nonce that
1652
                // was sent in our initial INFO and must be signed in CONNECT.
1653
                if didSolicit {
1,117✔
1654
                        c.nonce = []byte(info.Nonce)
548✔
1655
                }
548✔
1656
                if info.TLSRequired && didSolicit {
599✔
1657
                        remote.TLS = true
30✔
1658
                }
30✔
1659
                supportsHeaders := c.srv.supportsHeaders()
569✔
1660
                c.headers = supportsHeaders && info.Headers
569✔
1661

569✔
1662
                // Remember the remote server.
569✔
1663
                // Pre 2.2.0 servers are not sending their server name.
569✔
1664
                // In that case, use info.ID, which, for those servers, matches
569✔
1665
                // the content of the field `Name` in the leafnode CONNECT protocol.
569✔
1666
                if info.Name == _EMPTY_ {
569✔
1667
                        c.leaf.remoteServer = info.ID
×
1668
                } else {
569✔
1669
                        c.leaf.remoteServer = info.Name
569✔
1670
                }
569✔
1671
                c.leaf.remoteDomain = info.Domain
569✔
1672
                c.leaf.remoteCluster = info.Cluster
569✔
1673
                // We send the protocol version in the INFO protocol.
569✔
1674
                // Keep track of it, so we know if this connection supports message
569✔
1675
                // tracing for instance.
569✔
1676
                c.opts.Protocol = info.Proto
569✔
1677
        }
1678

1679
        // For both initial INFO and async INFO protocols, Possibly
1680
        // update our list of remote leafnode URLs we can connect to,
1681
        // unless we are instructed not to.
1682
        if didSolicit && !remote.IgnoreDiscoveredServers &&
1,117✔
1683
                (len(info.LeafNodeURLs) > 0 || len(info.WSConnectURLs) > 0) {
2,184✔
1684
                // Consider the incoming array as the most up-to-date
1,067✔
1685
                // representation of the remote cluster's list of URLs.
1,067✔
1686
                c.updateLeafNodeURLs(info)
1,067✔
1687
        }
1,067✔
1688

1689
        // Only solicited leafnode connections trust permission updates from INFO.
1690
        if didSolicit && (info.Import != nil || info.Export != nil) {
1,123✔
1691
                perms := &Permissions{
6✔
1692
                        Publish:   info.Export,
6✔
1693
                        Subscribe: info.Import,
6✔
1694
                }
6✔
1695
                // Check if we have local deny clauses that we need to merge.
6✔
1696
                if remote := c.leaf.remote; remote != nil {
12✔
1697
                        if len(remote.DenyExports) > 0 {
7✔
1698
                                if perms.Publish == nil {
1✔
1699
                                        perms.Publish = &SubjectPermission{}
×
1700
                                }
×
1701
                                perms.Publish.Deny = append(perms.Publish.Deny, remote.DenyExports...)
1✔
1702
                        }
1703
                        if len(remote.DenyImports) > 0 {
7✔
1704
                                if perms.Subscribe == nil {
1✔
1705
                                        perms.Subscribe = &SubjectPermission{}
×
1706
                                }
×
1707
                                perms.Subscribe.Deny = append(perms.Subscribe.Deny, remote.DenyImports...)
1✔
1708
                        }
1709
                }
1710
                c.setPermissions(perms)
6✔
1711
        }
1712

1713
        var resumeConnect bool
1,117✔
1714

1,117✔
1715
        // If this is a remote connection and this is the first INFO protocol,
1,117✔
1716
        // then we need to finish the connect process by sending CONNECT, etc..
1,117✔
1717
        if firstINFO && didSolicit {
1,665✔
1718
                // Clear deadline that was set in createLeafNode while waiting for the INFO.
548✔
1719
                c.nc.SetDeadline(time.Time{})
548✔
1720
                resumeConnect = true
548✔
1721
        } else if !firstINFO && didSolicit {
1,637✔
1722
                c.leaf.remoteAccName = info.RemoteAccount
520✔
1723
        }
520✔
1724

1725
        // Check if we have the remote account information and if so make sure it's stored.
1726
        if info.RemoteAccount != _EMPTY_ {
1,628✔
1727
                if c.acc == nil {
511✔
1728
                        c.mu.Unlock()
×
1729
                        c.sendErr("Authorization Violation")
×
1730
                        c.closeConnection(ProtocolViolation)
×
1731
                        return
×
1732
                }
×
1733
                s.leafRemoteAccounts.Store(c.acc.Name, info.RemoteAccount)
511✔
1734
        }
1735
        c.mu.Unlock()
1,117✔
1736

1,117✔
1737
        finishConnect := info.ConnectInfo
1,117✔
1738
        if resumeConnect && s != nil {
1,665✔
1739
                s.leafNodeResumeConnectProcess(c)
548✔
1740
                if !info.InfoOnConnect {
548✔
1741
                        finishConnect = true
×
1742
                }
×
1743
        }
1744
        if finishConnect {
1,628✔
1745
                s.leafNodeFinishConnectProcess(c)
511✔
1746
        }
511✔
1747

1748
        // Check to see if we need to kick any internal source or mirror consumers.
1749
        // This will be a no-op if JetStream not enabled for this server or if the bound account
1750
        // does not have jetstream.
1751
        s.checkInternalSyncConsumers(c.acc)
1,117✔
1752
}
1753

1754
func (s *Server) negotiateLeafCompression(c *client, didSolicit bool, infoCompression string, co *CompressionOpts) (bool, error) {
992✔
1755
        // If WebSocket compression is already negotiated on this connection then
992✔
1756
        // we shouldn't layer S2 compression on top of it.
992✔
1757
        c.mu.Lock()
992✔
1758
        if c.ws != nil && c.ws.compress {
996✔
1759
                c.leaf.compression = CompressionOff
4✔
1760
                c.mu.Unlock()
4✔
1761
                return false, nil
4✔
1762
        }
4✔
1763
        c.mu.Unlock()
988✔
1764
        // Negotiate the appropriate compression mode (or no compression)
988✔
1765
        cm, err := selectCompressionMode(co.Mode, infoCompression)
988✔
1766
        if err != nil {
988✔
1767
                return false, err
×
1768
        }
×
1769
        c.mu.Lock()
988✔
1770
        // For "auto" mode, set the initial compression mode based on RTT
988✔
1771
        if cm == CompressionS2Auto {
1,862✔
1772
                if c.rttStart.IsZero() {
1,748✔
1773
                        c.rtt = computeRTT(c.start)
874✔
1774
                }
874✔
1775
                cm = selectS2AutoModeBasedOnRTT(c.rtt, co.RTTThresholds)
874✔
1776
        }
1777
        // Keep track of the negotiated compression mode.
1778
        c.leaf.compression = cm
988✔
1779
        cid := c.cid
988✔
1780
        var nonce string
988✔
1781
        if !didSolicit {
1,436✔
1782
                nonce = bytesToString(c.nonce)
448✔
1783
        }
448✔
1784
        c.mu.Unlock()
988✔
1785

988✔
1786
        if !needsCompression(cm) {
1,076✔
1787
                return false, nil
88✔
1788
        }
88✔
1789

1790
        // If we end-up doing compression...
1791

1792
        // Generate an INFO with the chosen compression mode.
1793
        s.mu.Lock()
900✔
1794
        info := s.copyLeafNodeInfo()
900✔
1795
        info.Compression, info.CID, info.Nonce = compressionModeForInfoProtocol(co, cm), cid, nonce
900✔
1796
        infoProto := generateInfoJSON(info)
900✔
1797
        s.mu.Unlock()
900✔
1798

900✔
1799
        // If we solicited, then send this INFO protocol BEFORE switching
900✔
1800
        // to compression writer. However, if we did not, we send it after.
900✔
1801
        c.mu.Lock()
900✔
1802
        if didSolicit {
1,352✔
1803
                c.enqueueProto(infoProto)
452✔
1804
                // Make sure it is completely flushed (the pending bytes goes to
452✔
1805
                // 0) before proceeding.
452✔
1806
                for c.out.pb > 0 && !c.isClosed() {
903✔
1807
                        c.flushOutbound()
451✔
1808
                }
451✔
1809
        }
1810
        // This is to notify the readLoop that it should switch to a
1811
        // (de)compression reader.
1812
        c.in.flags.set(switchToCompression)
900✔
1813
        // Create the compress writer before queueing the INFO protocol for
900✔
1814
        // a route that did not solicit. It will make sure that that proto
900✔
1815
        // is sent with compression on.
900✔
1816
        c.out.cw = s2.NewWriter(nil, s2WriterOptions(cm)...)
900✔
1817
        if !didSolicit {
1,348✔
1818
                c.enqueueProto(infoProto)
448✔
1819
        }
448✔
1820
        c.mu.Unlock()
900✔
1821
        return true, nil
900✔
1822
}
1823

1824
// When getting a leaf node INFO protocol, use the provided
1825
// array of urls to update the list of possible endpoints.
1826
func (c *client) updateLeafNodeURLs(info *Info) {
1,067✔
1827
        cfg := c.leaf.remote
1,067✔
1828
        cfg.Lock()
1,067✔
1829
        defer cfg.Unlock()
1,067✔
1830

1,067✔
1831
        // We have ensured that if a remote has a WS scheme, then all are.
1,067✔
1832
        // So check if first is WS, then add WS URLs, otherwise, add non WS ones.
1,067✔
1833
        if len(cfg.URLs) > 0 && isWSURL(cfg.URLs[0]) {
1,119✔
1834
                // It does not really matter if we use "ws://" or "wss://" here since
52✔
1835
                // we will have already marked that the remote should use TLS anyway.
52✔
1836
                // But use proper scheme for log statements, etc...
52✔
1837
                proto := wsSchemePrefix
52✔
1838
                if cfg.TLS {
52✔
1839
                        proto = wsSchemePrefixTLS
×
1840
                }
×
1841
                c.doUpdateLNURLs(cfg, proto, info.WSConnectURLs)
52✔
1842
                return
52✔
1843
        }
1844
        c.doUpdateLNURLs(cfg, "nats-leaf", info.LeafNodeURLs)
1,015✔
1845
}
1846

1847
func (c *client) doUpdateLNURLs(cfg *leafNodeCfg, scheme string, URLs []string) {
1,067✔
1848
        cfg.urls = make([]*url.URL, 0, 1+len(URLs))
1,067✔
1849
        // Add the ones we receive in the protocol
1,067✔
1850
        for _, surl := range URLs {
2,916✔
1851
                url, err := url.Parse(fmt.Sprintf("%s://%s", scheme, surl))
1,849✔
1852
                if err != nil {
1,849✔
1853
                        // As per below, the URLs we receive should not have contained URL info, so this should be safe to log.
×
1854
                        c.Errorf("Error parsing url %q: %v", surl, err)
×
1855
                        continue
×
1856
                }
1857
                // Do not add if it's the same as what we already have configured.
1858
                var dup bool
1,849✔
1859
                for _, u := range cfg.URLs {
4,571✔
1860
                        // URLs that we receive never have user info, but the
2,722✔
1861
                        // ones that were configured may have. Simply compare
2,722✔
1862
                        // host and port to decide if they are equal or not.
2,722✔
1863
                        if url.Host == u.Host && url.Port() == u.Port() {
4,098✔
1864
                                dup = true
1,376✔
1865
                                break
1,376✔
1866
                        }
1867
                }
1868
                if !dup {
2,322✔
1869
                        cfg.urls = append(cfg.urls, url)
473✔
1870
                        cfg.saveTLSHostname(url)
473✔
1871
                }
473✔
1872
        }
1873
        // Add the configured one
1874
        cfg.urls = append(cfg.urls, cfg.URLs...)
1,067✔
1875
}
1876

1877
// Similar to setInfoHostPortAndGenerateJSON, but for leafNodeInfo.
1878
func (s *Server) setLeafNodeInfoHostPortAndIP() error {
3,982✔
1879
        opts := s.getOpts()
3,982✔
1880
        if opts.LeafNode.Advertise != _EMPTY_ {
3,993✔
1881
                advHost, advPort, err := parseHostPort(opts.LeafNode.Advertise, opts.LeafNode.Port)
11✔
1882
                if err != nil {
11✔
1883
                        return err
×
1884
                }
×
1885
                s.leafNodeInfo.Host = advHost
11✔
1886
                s.leafNodeInfo.Port = advPort
11✔
1887
        } else {
3,971✔
1888
                s.leafNodeInfo.Host = opts.LeafNode.Host
3,971✔
1889
                s.leafNodeInfo.Port = opts.LeafNode.Port
3,971✔
1890
                // If the host is "0.0.0.0" or "::" we need to resolve to a public IP.
3,971✔
1891
                // This will return at most 1 IP.
3,971✔
1892
                hostIsIPAny, ips, err := s.getNonLocalIPsIfHostIsIPAny(s.leafNodeInfo.Host, false)
3,971✔
1893
                if err != nil {
3,971✔
1894
                        return err
×
1895
                }
×
1896
                if hostIsIPAny {
4,204✔
1897
                        if len(ips) == 0 {
233✔
1898
                                s.Errorf("Could not find any non-local IP for leafnode's listen specification %q",
×
1899
                                        s.leafNodeInfo.Host)
×
1900
                        } else {
233✔
1901
                                // Take the first from the list...
233✔
1902
                                s.leafNodeInfo.Host = ips[0]
233✔
1903
                        }
233✔
1904
                }
1905
        }
1906
        // Use just host:port for the IP
1907
        s.leafNodeInfo.IP = net.JoinHostPort(s.leafNodeInfo.Host, strconv.Itoa(s.leafNodeInfo.Port))
3,982✔
1908
        if opts.LeafNode.Advertise != _EMPTY_ {
3,993✔
1909
                s.Noticef("Advertise address for leafnode is set to %s", s.leafNodeInfo.IP)
11✔
1910
        }
11✔
1911
        return nil
3,982✔
1912
}
1913

1914
// Add the connection to the map of leaf nodes.
1915
// If `checkForDup` is true (invoked when a leafnode is accepted), then we check
1916
// if a connection already exists for the same server name and account.
1917
// That can happen when the remote is attempting to reconnect while the accepting
1918
// side did not detect the connection as broken yet.
1919
// But it can also happen when there is a misconfiguration and the remote is
1920
// creating two (or more) connections that bind to the same account on the accept
1921
// side.
1922
// When a duplicate is found, the new connection is accepted and the old is closed
1923
// (this solves the stale connection situation). An error is returned to help the
1924
// remote detect the misconfiguration when the duplicate is the result of that
1925
// misconfiguration.
1926
func (s *Server) addLeafNodeConnection(c *client, srvName, clusterName string, checkForDup bool) bool {
1,058✔
1927
        var accName string
1,058✔
1928
        c.mu.Lock()
1,058✔
1929
        cid := c.cid
1,058✔
1930
        acc := c.acc
1,058✔
1931
        if acc != nil {
2,116✔
1932
                accName = acc.Name
1,058✔
1933
        }
1,058✔
1934
        myRemoteDomain := c.leaf.remoteDomain
1,058✔
1935
        mySrvName := c.leaf.remoteServer
1,058✔
1936
        remoteAccName := c.leaf.remoteAccName
1,058✔
1937
        myClustName := c.leaf.remoteCluster
1,058✔
1938
        remote := c.leaf.remote
1,058✔
1939
        solicited := remote != nil
1,058✔
1940
        c.mu.Unlock()
1,058✔
1941

1,058✔
1942
        var old *client
1,058✔
1943
        s.mu.Lock()
1,058✔
1944
        // We check for empty because in some test we may send empty CONNECT{}
1,058✔
1945
        if checkForDup && srvName != _EMPTY_ {
1,569✔
1946
                for _, ol := range s.leafs {
840✔
1947
                        ol.mu.Lock()
329✔
1948
                        // We care here only about non solicited Leafnode. This function
329✔
1949
                        // is more about replacing stale connections than detecting loops.
329✔
1950
                        // We have code for the loop detection elsewhere, which also delays
329✔
1951
                        // attempt to reconnect.
329✔
1952
                        if !ol.isSolicitedLeafNode() && ol.leaf.remoteServer == srvName &&
329✔
1953
                                ol.leaf.remoteCluster == clusterName && ol.acc.Name == accName &&
329✔
1954
                                remoteAccName != _EMPTY_ && ol.leaf.remoteAccName == remoteAccName {
331✔
1955
                                old = ol
2✔
1956
                        }
2✔
1957
                        ol.mu.Unlock()
329✔
1958
                        if old != nil {
331✔
1959
                                break
2✔
1960
                        }
1961
                }
1962
        }
1963
        // Now that we are under the server lock and before adding it to the map,
1964
        // for a solicited leaf, we need to make sure that it has not been removed
1965
        // from the config or disabled.
1966
        if solicited {
1,567✔
1967
                // If no longer valid, do not add to the server map. The connection
509✔
1968
                // should have been marked so that it can't reconnect. When the caller
509✔
1969
                // calls closeConnection(), cleanup (including clearing the connect-
509✔
1970
                // in-progress flag) will occur at the appropriate time.
509✔
1971
                if !remote.stillValid() {
509✔
1972
                        // Prevent reconnect in case it was not yet done.
×
1973
                        c.setNoReconnect()
×
1974
                        s.mu.Unlock()
×
1975
                        s.removeFromTempClients(cid)
×
1976
                        return false
×
1977
                }
×
1978
                remote.setConnectInProgress(false)
509✔
1979
        }
1980
        // Store new connection in the map
1981
        s.leafs[cid] = c
1,058✔
1982
        s.mu.Unlock()
1,058✔
1983
        s.removeFromTempClients(cid)
1,058✔
1984

1,058✔
1985
        // If applicable, evict the old one.
1,058✔
1986
        if old != nil {
1,060✔
1987
                old.sendErrAndErr(DuplicateRemoteLeafnodeConnection.String())
2✔
1988
                old.closeConnection(DuplicateRemoteLeafnodeConnection)
2✔
1989
                c.Warnf("Replacing connection from same server")
2✔
1990
        }
2✔
1991

1992
        srvDecorated := func() string {
1,249✔
1993
                if myClustName == _EMPTY_ {
213✔
1994
                        return mySrvName
22✔
1995
                }
22✔
1996
                return fmt.Sprintf("%s/%s", mySrvName, myClustName)
169✔
1997
        }
1998

1999
        opts := s.getOpts()
1,058✔
2000
        sysAcc := s.SystemAccount()
1,058✔
2001
        js := s.getJetStream()
1,058✔
2002
        var meta *raft
1,058✔
2003
        if js != nil {
1,517✔
2004
                if mg := js.getMetaGroup(); mg != nil {
802✔
2005
                        meta = mg.(*raft)
343✔
2006
                }
343✔
2007
        }
2008
        blockMappingOutgoing := false
1,058✔
2009
        // Deny (non domain) JetStream API traffic unless system account is shared
1,058✔
2010
        // and domain names are identical and extending is not disabled
1,058✔
2011

1,058✔
2012
        // Check if backwards compatibility has been enabled and needs to be acted on
1,058✔
2013
        forceSysAccDeny := false
1,058✔
2014
        if len(opts.JsAccDefaultDomain) > 0 {
1,091✔
2015
                if acc == sysAcc {
44✔
2016
                        for _, d := range opts.JsAccDefaultDomain {
22✔
2017
                                if d == _EMPTY_ {
19✔
2018
                                        // Extending JetStream via leaf node is mutually exclusive with a domain mapping to the empty/default domain.
8✔
2019
                                        // As soon as one mapping to "" is found, disable the ability to extend JS via a leaf node.
8✔
2020
                                        c.Noticef("Not extending remote JetStream domain %q due to presence of empty default domain", myRemoteDomain)
8✔
2021
                                        forceSysAccDeny = true
8✔
2022
                                        break
8✔
2023
                                }
2024
                        }
2025
                } else if domain, ok := opts.JsAccDefaultDomain[accName]; ok && domain == _EMPTY_ {
35✔
2026
                        // for backwards compatibility with old setups that do not have a domain name set
13✔
2027
                        c.Debugf("Skipping deny %q for account %q due to default domain", jsAllAPI, accName)
13✔
2028
                        return true
13✔
2029
                }
13✔
2030
        }
2031

2032
        // If the server has JS disabled, it may still be part of a JetStream that could be extended.
2033
        // This is either signaled by js being disabled and a domain set,
2034
        // or in cases where no domain name exists, an extension hint is set.
2035
        // However, this is only relevant in mixed setups.
2036
        //
2037
        // If the system account connects but default domains are present, JetStream can't be extended.
2038
        if opts.JetStreamDomain != myRemoteDomain || (!opts.JetStream && (opts.JetStreamDomain == _EMPTY_ && opts.JetStreamExtHint != jsWillExtend)) ||
1,045✔
2039
                sysAcc == nil || acc == nil || forceSysAccDeny {
1,955✔
2040
                // If domain names mismatch always deny. This applies to system accounts as well as non system accounts.
910✔
2041
                // Not having a system account, account or JetStream disabled is considered a mismatch as well.
910✔
2042
                if acc != nil && acc == sysAcc {
1,035✔
2043
                        c.Noticef("System account connected from %s", srvDecorated())
125✔
2044
                        c.Noticef("JetStream not extended, domains differ")
125✔
2045
                        c.mergeDenyPermissionsLocked(both, denyAllJs)
125✔
2046
                        // When a remote with a system account is present in a server, unless otherwise disabled, the server will be
125✔
2047
                        // started in observer mode. Now that it is clear that this not used, turn the observer mode off.
125✔
2048
                        if solicited && meta != nil && meta.IsObserver() {
151✔
2049
                                meta.setObserver(false, extNotExtended)
26✔
2050
                                c.Debugf("Turning JetStream metadata controller Observer Mode off")
26✔
2051
                                // Take note that the domain was not extended to avoid this state from startup.
26✔
2052
                                writePeerState(c.srv.diskIOSemaphore(), js.config.StoreDir, meta.currentPeerState())
26✔
2053
                                // Meta controller can't be leader yet.
26✔
2054
                                // Yet it is possible that due to observer mode every server already stopped campaigning.
26✔
2055
                                // Therefore this server needs to be kicked into campaigning gear explicitly.
26✔
2056
                                meta.Campaign()
26✔
2057
                        }
26✔
2058
                } else {
785✔
2059
                        c.Noticef("JetStream using domains: local %q, remote %q", opts.JetStreamDomain, myRemoteDomain)
785✔
2060
                        c.mergeDenyPermissionsLocked(both, denyAllClientJs)
785✔
2061
                }
785✔
2062
                blockMappingOutgoing = true
910✔
2063
        } else if acc == sysAcc {
201✔
2064
                // system account and same domain
66✔
2065
                s.sys.client.Noticef("Extending JetStream domain %q as System Account connected from server %s",
66✔
2066
                        myRemoteDomain, srvDecorated())
66✔
2067
                // In an extension use case, pin leadership to server remotes connect to.
66✔
2068
                // Therefore, server with a remote that are not already in observer mode, need to be put into it.
66✔
2069
                if solicited && meta != nil && !meta.IsObserver() {
70✔
2070
                        c.Debugf("Turning JetStream metadata controller Observer Mode on - System Account Connected")
4✔
2071
                        // Discard any local metagroup state accumulated before the SYS-account
4✔
2072
                        // leaf came up (e.g. the wrong-hint case where this server bootstrapped
4✔
2073
                        // its own metagroup). The parent's view is now authoritative; without
4✔
2074
                        // this reset the two raft logs stay forked because the standalone log's
4✔
2075
                        // commit prefix short-circuits the follower's AE handling.
4✔
2076
                        meta.setObserver(true, extExtended)
4✔
2077
                        meta.Reset()
4✔
2078
                }
4✔
2079
        } else {
69✔
2080
                // This deny is needed in all cases (system account shared or not)
69✔
2081
                // If the system account is shared, jsAllAPI traffic will go through the system account.
69✔
2082
                // So in order to prevent duplicate delivery (from system and actual account) suppress it on the account.
69✔
2083
                // If the system account is NOT shared, jsAllAPI traffic has no business
69✔
2084
                c.Debugf("Adding deny %+v for account %q", denyAllClientJs, accName)
69✔
2085
                c.mergeDenyPermissionsLocked(both, denyAllClientJs)
69✔
2086
        }
69✔
2087
        // If we have a specified JetStream domain we will want to add a mapping to
2088
        // allow access cross domain for each non-system account.
2089
        if opts.JetStreamDomain != _EMPTY_ && opts.JetStream && acc != nil && acc != sysAcc {
1,273✔
2090
                for src, dest := range generateJSMappingTable(opts.JetStreamDomain) {
2,280✔
2091
                        if err := acc.AddMapping(src, dest); err != nil {
2,052✔
2092
                                c.Debugf("Error adding JetStream domain mapping: %s", err.Error())
×
2093
                        } else {
2,052✔
2094
                                c.Debugf("Adding JetStream Domain Mapping %q -> %s to account %q", src, dest, accName)
2,052✔
2095
                        }
2,052✔
2096
                }
2097
                if blockMappingOutgoing {
425✔
2098
                        src := fmt.Sprintf(jsDomainAPI, opts.JetStreamDomain)
197✔
2099
                        // make sure that messages intended for this domain, do not leave the cluster via this leaf node connection
197✔
2100
                        // This is a guard against a miss-config with two identical domain names and will only cover some forms
197✔
2101
                        // of this issue, not all of them.
197✔
2102
                        // This guards against a hub and a spoke having the same domain name.
197✔
2103
                        // But not two spokes having the same one and the request coming from the hub.
197✔
2104
                        c.mergeDenyPermissionsLocked(pub, []string{src})
197✔
2105
                        c.Debugf("Adding deny %q for outgoing messages to account %q", src, accName)
197✔
2106
                }
197✔
2107
        }
2108
        return true
1,045✔
2109
}
2110

2111
func (s *Server) removeLeafNodeConnection(c *client) {
1,405✔
2112
        s.mu.Lock()
1,405✔
2113
        c.mu.Lock()
1,405✔
2114
        cid := c.cid
1,405✔
2115
        if c.leaf != nil {
2,809✔
2116
                if c.leaf.tsubt != nil {
2,362✔
2117
                        c.leaf.tsubt.Stop()
958✔
2118
                        c.leaf.tsubt = nil
958✔
2119
                }
958✔
2120
                if c.leaf.gwSub != nil {
1,912✔
2121
                        s.gwLeafSubs.Remove(c.leaf.gwSub)
508✔
2122
                        // We need to set this to nil for GC to release the connection
508✔
2123
                        c.leaf.gwSub = nil
508✔
2124
                }
508✔
2125
                if remote := c.leaf.remote; remote != nil {
2,076✔
2126
                        // If "noReconnect" is true, then we won't attempt to reconnect, so
672✔
2127
                        // we will clear the "connect-in-progress" flag. However, if we can
672✔
2128
                        // reconnect, then we should set "connect-in-progress" to true while
672✔
2129
                        // we are under the server/client lock. The go routine that performs
672✔
2130
                        // the reconnect will be started later and there would be a gap with
672✔
2131
                        // the wrong flag value otherwise.
672✔
2132
                        remote.setConnectInProgress(!c.flags.isSet(noReconnect))
672✔
2133
                }
672✔
2134
        }
2135
        proxyKey := c.proxyKey
1,405✔
2136
        c.mu.Unlock()
1,405✔
2137
        delete(s.leafs, cid)
1,405✔
2138
        if proxyKey != _EMPTY_ {
1,409✔
2139
                s.removeProxiedConn(proxyKey, cid)
4✔
2140
        }
4✔
2141
        s.mu.Unlock()
1,405✔
2142
        s.removeFromTempClients(cid)
1,405✔
2143
}
2144

2145
// Connect information for solicited leafnodes.
2146
type leafConnectInfo struct {
2147
        Version   string   `json:"version,omitempty"`
2148
        Nkey      string   `json:"nkey,omitempty"`
2149
        JWT       string   `json:"jwt,omitempty"`
2150
        Sig       string   `json:"sig,omitempty"`
2151
        User      string   `json:"user,omitempty"`
2152
        Pass      string   `json:"pass,omitempty"`
2153
        Token     string   `json:"auth_token,omitempty"`
2154
        ID        string   `json:"server_id,omitempty"`
2155
        Domain    string   `json:"domain,omitempty"`
2156
        Name      string   `json:"name,omitempty"`
2157
        Hub       bool     `json:"is_hub,omitempty"`
2158
        Cluster   string   `json:"cluster,omitempty"`
2159
        Headers   bool     `json:"headers,omitempty"`
2160
        JetStream bool     `json:"jetstream,omitempty"`
2161
        DenyPub   []string `json:"deny_pub,omitempty"`
2162
        Isolate   bool     `json:"isolate,omitempty"`
2163

2164
        // There was an existing field called:
2165
        // >> Comp bool `json:"compression,omitempty"`
2166
        // that has never been used. With support for compression, we now need
2167
        // a field that is a string. So we use a different json tag:
2168
        Compression string `json:"compress_mode,omitempty"`
2169

2170
        // Just used to detect wrong connection attempts.
2171
        Gateway string `json:"gateway,omitempty"`
2172

2173
        // Tells the accept side which account the remote is binding to.
2174
        RemoteAccount string `json:"remote_account,omitempty"`
2175

2176
        // The accept side of a LEAF connection, unlike ROUTER and GATEWAY, receives
2177
        // only the CONNECT protocol, and no INFO. So we need to send the protocol
2178
        // version as part of the CONNECT. It will indicate if a connection supports
2179
        // some features, such as message tracing.
2180
        // We use `protocol` as the JSON tag, so this is automatically unmarshal'ed
2181
        // in the low level process CONNECT.
2182
        Proto int `json:"protocol,omitempty"`
2183
}
2184

2185
// processLeafNodeConnect will process the inbound connect args.
2186
// Once we are here we are bound to an account, so can send any interest that
2187
// we would have to the other side.
2188
func (c *client) processLeafNodeConnect(s *Server, arg []byte, lang string) error {
554✔
2189
        // Way to detect clients that incorrectly connect to the route listen
554✔
2190
        // port. Client provided "lang" in the CONNECT protocol while LEAFNODEs don't.
554✔
2191
        if lang != _EMPTY_ {
554✔
2192
                c.sendErrAndErr(ErrClientConnectedToLeafNodePort.Error())
×
2193
                c.closeConnection(WrongPort)
×
2194
                return ErrClientConnectedToLeafNodePort
×
2195
        }
×
2196

2197
        // Unmarshal as a leaf node connect protocol
2198
        proto := &leafConnectInfo{}
554✔
2199
        if err := json.Unmarshal(arg, proto); err != nil {
554✔
2200
                return err
×
2201
        }
×
2202

2203
        // Reject a cluster that contains spaces.
2204
        if proto.Cluster != _EMPTY_ && strings.Contains(proto.Cluster, " ") {
555✔
2205
                c.sendErrAndErr(ErrClusterNameHasSpaces.Error())
1✔
2206
                c.closeConnection(ProtocolViolation)
1✔
2207
                return ErrClusterNameHasSpaces
1✔
2208
        }
1✔
2209

2210
        // Check for cluster name collisions.
2211
        if cn := s.cachedClusterName(); cn != _EMPTY_ && proto.Cluster != _EMPTY_ && proto.Cluster == cn {
556✔
2212
                c.sendErrAndErr(ErrLeafNodeHasSameClusterName.Error())
3✔
2213
                c.closeConnection(ClusterNamesIdentical)
3✔
2214
                return ErrLeafNodeHasSameClusterName
3✔
2215
        }
3✔
2216

2217
        // Reject if this has Gateway which means that it would be from a gateway
2218
        // connection that incorrectly connects to the leafnode port.
2219
        if proto.Gateway != _EMPTY_ {
550✔
2220
                errTxt := fmt.Sprintf("Rejecting connection from gateway %q on the leafnode port", proto.Gateway)
×
2221
                c.Errorf(errTxt)
×
2222
                c.sendErr(errTxt)
×
2223
                c.closeConnection(WrongGateway)
×
2224
                return ErrWrongGateway
×
2225
        }
×
2226

2227
        if mv := s.getOpts().LeafNode.MinVersion; mv != _EMPTY_ {
552✔
2228
                major, minor, update, _ := versionComponents(mv)
2✔
2229
                if !versionAtLeast(proto.Version, major, minor, update) {
3✔
2230
                        // Send back an INFO so recent remote servers process the rejection
1✔
2231
                        // cleanly, then close immediately. The soliciting side applies the
1✔
2232
                        // reconnect delay when it processes the error.
1✔
2233
                        s.sendPermsAndAccountInfo(c)
1✔
2234
                        c.sendErrAndErr(fmt.Sprintf("%s %q", ErrLeafNodeMinVersionRejected, mv))
1✔
2235
                        c.closeConnection(MinimumVersionRequired)
1✔
2236
                        return ErrMinimumVersionRequired
1✔
2237
                }
1✔
2238
        }
2239

2240
        // Check if this server supports headers.
2241
        supportHeaders := c.srv.supportsHeaders()
549✔
2242

549✔
2243
        c.mu.Lock()
549✔
2244
        // Leaf Nodes do not do echo or verbose or pedantic.
549✔
2245
        c.opts.Verbose = false
549✔
2246
        c.opts.Echo = false
549✔
2247
        c.opts.Pedantic = false
549✔
2248
        // This inbound connection will be marked as supporting headers if this server
549✔
2249
        // support headers and the remote has sent in the CONNECT protocol that it does
549✔
2250
        // support headers too.
549✔
2251
        c.headers = supportHeaders && proto.Headers
549✔
2252
        // If the compression level is still not set, set it based on what has been
549✔
2253
        // given to us in the CONNECT protocol.
549✔
2254
        if c.leaf.compression == _EMPTY_ {
679✔
2255
                // But if proto.Compression is _EMPTY_, set it to CompressionNotSupported
130✔
2256
                if proto.Compression == _EMPTY_ {
168✔
2257
                        c.leaf.compression = CompressionNotSupported
38✔
2258
                } else {
130✔
2259
                        c.leaf.compression = proto.Compression
92✔
2260
                }
92✔
2261
        }
2262

2263
        // Remember the remote server.
2264
        c.leaf.remoteServer = proto.Name
549✔
2265
        // Remember the remote account name
549✔
2266
        c.leaf.remoteAccName = proto.RemoteAccount
549✔
2267
        // Remember if the leafnode requested isolation.
549✔
2268
        c.leaf.isolated = c.leaf.isolated || proto.Isolate
549✔
2269

549✔
2270
        // If the other side has declared itself a hub, so we will take on the spoke role.
549✔
2271
        if proto.Hub {
561✔
2272
                c.leaf.isSpoke = true
12✔
2273
        }
12✔
2274

2275
        // The soliciting side is part of a cluster.
2276
        if proto.Cluster != _EMPTY_ {
953✔
2277
                c.leaf.remoteCluster = proto.Cluster
404✔
2278
        }
404✔
2279

2280
        c.leaf.remoteDomain = proto.Domain
549✔
2281

549✔
2282
        // When a leaf solicits a connection to a hub, the perms that it will use on the soliciting leafnode's
549✔
2283
        // behalf are correct for them, but inside the hub need to be reversed since data is flowing in the opposite direction.
549✔
2284
        if !c.isSolicitedLeafNode() && c.perms != nil {
555✔
2285
                sp, pp := c.perms.sub, c.perms.pub
6✔
2286
                c.perms.sub, c.perms.pub = pp, sp
6✔
2287
                if c.opts.Import != nil {
11✔
2288
                        c.darray = c.opts.Import.Deny
5✔
2289
                } else {
6✔
2290
                        c.darray = nil
1✔
2291
                }
1✔
2292
        }
2293

2294
        // Set the Ping timer
2295
        c.setFirstPingTimer()
549✔
2296

549✔
2297
        // If we received pub deny permissions from the other end, merge with existing ones.
549✔
2298
        c.mergeDenyPermissions(pub, proto.DenyPub)
549✔
2299

549✔
2300
        acc := c.acc
549✔
2301
        c.mu.Unlock()
549✔
2302

549✔
2303
        // If the account is not set (e.g. connection was closed due to auth
549✔
2304
        // timeout while still being processed), bail out to avoid a panic.
549✔
2305
        if acc == nil {
549✔
2306
                c.closeConnection(MissingAccount)
×
2307
                return ErrMissingAccount
×
2308
        }
×
2309

2310
        // Register the cluster, even if empty, as long as we are acting as a hub.
2311
        if !proto.Hub {
1,086✔
2312
                acc.registerLeafNodeCluster(proto.Cluster)
537✔
2313
        }
537✔
2314

2315
        // Add in the leafnode here since we passed through auth at this point.
2316
        s.addLeafNodeConnection(c, proto.Name, proto.Cluster, true)
549✔
2317

549✔
2318
        // If we have permissions bound to this leafnode we need to send then back to the
549✔
2319
        // origin server for local enforcement.
549✔
2320
        s.sendPermsAndAccountInfo(c)
549✔
2321

549✔
2322
        // Create and initialize the smap since we know our bound account now.
549✔
2323
        // This will send all registered subs too.
549✔
2324
        s.initLeafNodeSmapAndSendSubs(c)
549✔
2325

549✔
2326
        // Announce the account connect event for a leaf node.
549✔
2327
        // This will be a no-op as needed.
549✔
2328
        s.sendLeafNodeConnect(c.acc)
549✔
2329

549✔
2330
        // Check to see if we need to kick any internal source or mirror consumers.
549✔
2331
        // This will be a no-op if JetStream not enabled for this server or if the bound account
549✔
2332
        // does not have jetstream.
549✔
2333
        s.checkInternalSyncConsumers(acc)
549✔
2334

549✔
2335
        return nil
549✔
2336
}
2337

2338
// checkInternalSyncConsumers
2339
func (s *Server) checkInternalSyncConsumers(acc *Account) {
1,666✔
2340
        // Grab our js
1,666✔
2341
        js := s.getJetStream()
1,666✔
2342

1,666✔
2343
        // Only applicable if we have JS and the leafnode has JS as well.
1,666✔
2344
        // We check for remote JS outside.
1,666✔
2345
        if !js.isEnabled() || acc == nil {
2,592✔
2346
                return
926✔
2347
        }
926✔
2348

2349
        // We will check all streams in our local account. They must be a leader and
2350
        // be sourcing or mirroring. We will check the external config on the stream itself
2351
        // if this is cross domain, or if the remote domain is empty, meaning we might be
2352
        // extending the system across this leafnode connection and hence we would be extending
2353
        // our own domain.
2354
        jsa := js.lookupAccount(acc)
740✔
2355
        if jsa == nil {
1,017✔
2356
                return
277✔
2357
        }
277✔
2358

2359
        var streams []*stream
463✔
2360
        jsa.mu.RLock()
463✔
2361
        for _, mset := range jsa.streams {
528✔
2362
                mset.cfgMu.RLock()
65✔
2363
                // We need to have a mirror or source defined.
65✔
2364
                // We do not want to force another lock here to look for leader status,
65✔
2365
                // so collect and after we release jsa will make sure.
65✔
2366
                if mset.cfg.Mirror != nil || len(mset.cfg.Sources) > 0 {
78✔
2367
                        streams = append(streams, mset)
13✔
2368
                }
13✔
2369
                mset.cfgMu.RUnlock()
65✔
2370
        }
2371
        jsa.mu.RUnlock()
463✔
2372

463✔
2373
        // Now loop through all candidates and check if we are the leader and have NOT
463✔
2374
        // created the sync up consumer.
463✔
2375
        for _, mset := range streams {
476✔
2376
                mset.retryDisconnectedSyncConsumers()
13✔
2377
        }
13✔
2378
}
2379

2380
// Returns the remote cluster name. This is set only once so does not require a lock.
2381
func (c *client) remoteCluster() string {
118,864✔
2382
        if c.leaf == nil {
118,864✔
2383
                return _EMPTY_
×
2384
        }
×
2385
        return c.leaf.remoteCluster
118,864✔
2386
}
2387

2388
// Sends back an info block to the soliciting leafnode to let it know about
2389
// its permission settings for local enforcement.
2390
func (s *Server) sendPermsAndAccountInfo(c *client) {
550✔
2391
        // Copy
550✔
2392
        s.mu.Lock()
550✔
2393
        info := s.copyLeafNodeInfo()
550✔
2394
        s.mu.Unlock()
550✔
2395
        c.mu.Lock()
550✔
2396
        info.CID = c.cid
550✔
2397
        info.Import = c.opts.Import
550✔
2398
        info.Export = c.opts.Export
550✔
2399
        info.RemoteAccount = c.acc.Name
550✔
2400
        // s.SystemAccount() uses an atomic operation and does not get the server lock, so this is safe.
550✔
2401
        info.IsSystemAccount = c.acc == s.SystemAccount()
550✔
2402
        info.ConnectInfo = true
550✔
2403
        c.enqueueProto(generateInfoJSON(info))
550✔
2404
        c.mu.Unlock()
550✔
2405
}
550✔
2406

2407
// Snapshot the current subscriptions from the sublist into our smap which
2408
// we will keep updated from now on.
2409
// Also send the registered subscriptions.
2410
func (s *Server) initLeafNodeSmapAndSendSubs(c *client) {
1,058✔
2411
        acc := c.acc
1,058✔
2412
        if acc == nil {
1,058✔
2413
                c.Debugf("Leafnode does not have an account bound")
×
2414
                return
×
2415
        }
×
2416
        // Collect all account subs here.
2417
        _subs := [1024]*subscription{}
1,058✔
2418
        subs := _subs[:0]
1,058✔
2419
        ims := []string{}
1,058✔
2420

1,058✔
2421
        // Hold the client lock otherwise there can be a race and miss some subs.
1,058✔
2422
        c.mu.Lock()
1,058✔
2423
        defer c.mu.Unlock()
1,058✔
2424

1,058✔
2425
        acc.mu.RLock()
1,058✔
2426
        accName := acc.Name
1,058✔
2427
        accNTag := acc.nameTag
1,058✔
2428

1,058✔
2429
        // To make printing look better when no friendly name present.
1,058✔
2430
        if accNTag != _EMPTY_ {
1,061✔
2431
                accNTag = "/" + accNTag
3✔
2432
        }
3✔
2433

2434
        // If we are solicited we only send interest for local clients.
2435
        if c.isSpokeLeafNode() {
1,567✔
2436
                acc.sl.localSubs(&subs, true)
509✔
2437
        } else {
1,058✔
2438
                acc.sl.All(&subs)
549✔
2439
        }
549✔
2440

2441
        // Check if we have an existing service import reply.
2442
        siReply := copyBytes(acc.siReply)
1,058✔
2443

1,058✔
2444
        // Since leaf nodes only send on interest, if the bound
1,058✔
2445
        // account has import services we need to send those over.
1,058✔
2446
        for isubj := range acc.imports.services {
4,920✔
2447
                if c.isSpokeLeafNode() && !c.canSubscribe(isubj) {
4,095✔
2448
                        c.Debugf("Not permitted to import service %q on behalf of %s%s", isubj, accName, accNTag)
233✔
2449
                        continue
233✔
2450
                }
2451
                ims = append(ims, isubj)
3,629✔
2452
        }
2453
        // Likewise for mappings.
2454
        for _, m := range acc.mappings {
3,200✔
2455
                if c.isSpokeLeafNode() && !c.canSubscribe(m.src) {
2,160✔
2456
                        c.Debugf("Not permitted to import mapping %q on behalf of %s%s", m.src, accName, accNTag)
18✔
2457
                        continue
18✔
2458
                }
2459
                ims = append(ims, m.src)
2,124✔
2460
        }
2461

2462
        // Create a unique subject that will be used for loop detection.
2463
        lds := acc.lds
1,058✔
2464
        acc.mu.RUnlock()
1,058✔
2465

1,058✔
2466
        // Check if we have to create the LDS.
1,058✔
2467
        if lds == _EMPTY_ {
1,884✔
2468
                lds = leafNodeLoopDetectionSubjectPrefix + nuid.Next()
826✔
2469
                acc.mu.Lock()
826✔
2470
                acc.lds = lds
826✔
2471
                acc.mu.Unlock()
826✔
2472
        }
826✔
2473

2474
        // Now check for gateway interest. Leafnodes will put this into
2475
        // the proper mode to propagate, but they are not held in the account.
2476
        gwsa := [16]*client{}
1,058✔
2477
        gws := gwsa[:0]
1,058✔
2478
        s.getOutboundGatewayConnections(&gws)
1,058✔
2479
        for _, cgw := range gws {
1,117✔
2480
                cgw.mu.Lock()
59✔
2481
                gw := cgw.gw
59✔
2482
                cgw.mu.Unlock()
59✔
2483
                if gw != nil {
118✔
2484
                        if ei, _ := gw.outsim.Load(accName); ei != nil {
118✔
2485
                                if e := ei.(*outsie); e != nil && e.sl != nil {
118✔
2486
                                        e.sl.All(&subs)
59✔
2487
                                }
59✔
2488
                        }
2489
                }
2490
        }
2491

2492
        applyGlobalRouting := s.gateway.enabled
1,058✔
2493
        if c.isSpokeLeafNode() {
1,567✔
2494
                // Add a fake subscription for this solicited leafnode connection
509✔
2495
                // so that we can send back directly for mapped GW replies.
509✔
2496
                // We need to keep track of this subscription so it can be removed
509✔
2497
                // when the connection is closed so that the GC can release it.
509✔
2498
                c.leaf.gwSub = &subscription{client: c, subject: []byte(gwReplyPrefix + ">")}
509✔
2499
                c.srv.gwLeafSubs.Insert(c.leaf.gwSub)
509✔
2500
        }
509✔
2501

2502
        // Now walk the results and add them to our smap
2503
        rc := c.leaf.remoteCluster
1,058✔
2504
        c.leaf.smap = make(map[string]int32)
1,058✔
2505
        for _, sub := range subs {
35,033✔
2506
                // Check perms regardless of role.
33,975✔
2507
                if c.perms != nil && !c.canSubscribe(string(sub.subject)) {
36,087✔
2508
                        c.Debugf("Not permitted to subscribe to %q on behalf of %s%s", sub.subject, accName, accNTag)
2,112✔
2509
                        continue
2,112✔
2510
                }
2511
                // Don't advertise interest from leafnodes to other isolated leafnodes.
2512
                if sub.client.kind == LEAF && c.isIsolatedLeafNode() {
31,863✔
2513
                        continue
×
2514
                }
2515
                // We ignore ourselves here.
2516
                // Also don't add the subscription if it has a origin cluster and the
2517
                // cluster name matches the one of the client we are sending to.
2518
                if c != sub.client && (sub.origin == nil || (bytesToString(sub.origin) != rc)) {
59,157✔
2519
                        count := int32(1)
27,294✔
2520
                        if len(sub.queue) > 0 && sub.qw > 0 {
27,305✔
2521
                                count = sub.qw
11✔
2522
                        }
11✔
2523
                        c.leaf.smap[keyFromSub(sub)] += count
27,294✔
2524
                        if c.leaf.tsub == nil {
28,288✔
2525
                                c.leaf.tsub = make(map[*subscription]struct{})
994✔
2526
                        }
994✔
2527
                        c.leaf.tsub[sub] = struct{}{}
27,294✔
2528
                }
2529
        }
2530
        // FIXME(dlc) - We need to update appropriately on an account claims update.
2531
        for _, isubj := range ims {
6,811✔
2532
                c.leaf.smap[isubj]++
5,753✔
2533
        }
5,753✔
2534
        // If we have gateways enabled we need to make sure the other side sends us responses
2535
        // that have been augmented from the original subscription.
2536
        // TODO(dlc) - Should we lock this down more?
2537
        if applyGlobalRouting {
1,130✔
2538
                c.leaf.smap[oldGWReplyPrefix+"*.>"]++
72✔
2539
                c.leaf.smap[gwReplyPrefix+">"]++
72✔
2540
        }
72✔
2541
        // Detect loops by subscribing to a specific subject and checking
2542
        // if this sub is coming back to us.
2543
        c.leaf.smap[lds]++
1,058✔
2544

1,058✔
2545
        // Check if we need to add an existing siReply to our map.
1,058✔
2546
        // This will be a prefix so add on the wildcard.
1,058✔
2547
        if siReply != nil {
1,076✔
2548
                wcsub := append(siReply, '>')
18✔
2549
                c.leaf.smap[string(wcsub)]++
18✔
2550
        }
18✔
2551
        // Queue all protocols. There is no max pending limit for LN connection,
2552
        // so we don't need chunking. The writes will happen from the writeLoop.
2553
        var b bytes.Buffer
1,058✔
2554
        for key, n := range c.leaf.smap {
24,687✔
2555
                c.writeLeafSub(&b, key, n)
23,629✔
2556
        }
23,629✔
2557
        if b.Len() > 0 {
2,116✔
2558
                c.enqueueProto(b.Bytes())
1,058✔
2559
        }
1,058✔
2560
        if c.leaf.tsub != nil {
2,053✔
2561
                // Clear the tsub map after 5 seconds.
995✔
2562
                c.leaf.tsubt = time.AfterFunc(5*time.Second, func() {
1,030✔
2563
                        c.mu.Lock()
35✔
2564
                        if c.leaf != nil {
70✔
2565
                                c.leaf.tsub = nil
35✔
2566
                                c.leaf.tsubt = nil
35✔
2567
                        }
35✔
2568
                        c.mu.Unlock()
35✔
2569
                })
2570
        }
2571
}
2572

2573
// updateInterestForAccountOnGateway called from gateway code when processing RS+ and RS-.
2574
func (s *Server) updateInterestForAccountOnGateway(accName string, sub *subscription, delta int32) {
194,145✔
2575
        // Since we're in the gateway's readLoop, and we would otherwise block, don't allow fetching.
194,145✔
2576
        acc, err := s.lookupOrFetchAccount(accName, false)
194,145✔
2577
        if acc == nil || err != nil {
194,419✔
2578
                s.Debugf("No or bad account for %q, failed to update interest from gateway", accName)
274✔
2579
                return
274✔
2580
        }
274✔
2581
        acc.updateLeafNodes(sub, delta)
193,871✔
2582
}
2583

2584
// updateLeafNodesEx will make sure to update the account smap for the subscription.
2585
// Will also forward to all leaf nodes as needed.
2586
// If `hubOnly` is true, then will update only leaf nodes that connect to this server
2587
// (that is, for which this server acts as a hub to them).
2588
func (acc *Account) updateLeafNodesEx(sub *subscription, delta int32, hubOnly bool) {
2,503,723✔
2589
        if acc == nil || sub == nil {
2,503,723✔
2590
                return
×
2591
        }
×
2592

2593
        // We will do checks for no leafnodes and same cluster here inline and under the
2594
        // general account read lock.
2595
        // If we feel we need to update the leafnodes we will do that out of line to avoid
2596
        // blocking routes or GWs.
2597

2598
        acc.mu.RLock()
2,503,723✔
2599
        // First check if we even have leafnodes here.
2,503,723✔
2600
        if acc.nleafs == 0 {
4,950,068✔
2601
                acc.mu.RUnlock()
2,446,345✔
2602
                return
2,446,345✔
2603
        }
2,446,345✔
2604

2605
        // Is this a loop detection subject.
2606
        isLDS := bytes.HasPrefix(sub.subject, []byte(leafNodeLoopDetectionSubjectPrefix))
57,378✔
2607

57,378✔
2608
        // Capture the cluster even if its empty.
57,378✔
2609
        var cluster string
57,378✔
2610
        if sub.origin != nil {
99,878✔
2611
                cluster = bytesToString(sub.origin)
42,500✔
2612
        }
42,500✔
2613

2614
        // If we have an isolated cluster we can return early, as long as it is not a loop detection subject.
2615
        // Empty clusters will return false for the check.
2616
        if !isLDS && acc.isLeafNodeClusterIsolated(cluster) {
74,910✔
2617
                acc.mu.RUnlock()
17,532✔
2618
                return
17,532✔
2619
        }
17,532✔
2620

2621
        // We can release the general account lock.
2622
        acc.mu.RUnlock()
39,846✔
2623

39,846✔
2624
        // We can hold the list lock here to avoid having to copy a large slice.
39,846✔
2625
        acc.lmu.RLock()
39,846✔
2626
        defer acc.lmu.RUnlock()
39,846✔
2627

39,846✔
2628
        // Do this once.
39,846✔
2629
        subject := string(sub.subject)
39,846✔
2630

39,846✔
2631
        // Walk the connected leafnodes from a random starting point to avoid
39,846✔
2632
        // concurrent callers all contending over leafs in the same order.
39,846✔
2633
        nleafs := len(acc.lleafs)
39,846✔
2634
        start := 0
39,846✔
2635
        if nleafs > 1 {
45,171✔
2636
                start = rand.Intn(nleafs)
5,325✔
2637
        }
5,325✔
2638
        for i := 0; i < nleafs; i++ {
87,683✔
2639
                ln := acc.lleafs[(start+i)%nleafs]
47,837✔
2640
                if ln == sub.client {
74,581✔
2641
                        continue
26,744✔
2642
                }
2643
                ln.mu.RLock()
21,093✔
2644
                // Don't advertise interest from leafnodes to other isolated leafnodes.
21,093✔
2645
                if sub.client.kind == LEAF && ln.isIsolatedLeafNode() {
21,093✔
2646
                        ln.mu.RUnlock()
×
2647
                        continue
×
2648
                }
2649
                // If `hubOnly` is true, it means that we want to update only leafnodes
2650
                // that connect to this server (so isHubLeafNode() would return `true`).
2651
                if hubOnly && !ln.isHubLeafNode() {
21,099✔
2652
                        ln.mu.RUnlock()
6✔
2653
                        continue
6✔
2654
                }
2655
                // Check to make sure this sub does not have an origin cluster that matches the leafnode.
2656
                // If skipped, make sure that we still let go the "$LDS." subscription that allows
2657
                // the detection of loops as long as different cluster.
2658
                clusterDifferent := cluster != ln.remoteCluster()
21,087✔
2659
                update := (isLDS && clusterDifferent) ||
21,087✔
2660
                        ((cluster == _EMPTY_ || clusterDifferent) && (delta <= 0 || ln.canSubscribeInternal(subject)))
21,087✔
2661
                ln.mu.RUnlock()
21,087✔
2662
                if update {
39,034✔
2663
                        ln.mu.Lock()
17,947✔
2664
                        // The leaf role, isolation mode, and remote cluster are stable
17,947✔
2665
                        // for the connection. Recheck canSubscribe here since permissions
17,947✔
2666
                        // can change, and to initializes mperms for wildcard subscriptions
17,947✔
2667
                        // that collide with deny rules.
17,947✔
2668
                        if isLDS || delta <= 0 || ln.canSubscribe(subject) {
35,894✔
2669
                                ln.updateSmap(sub, delta, isLDS)
17,947✔
2670
                        }
17,947✔
2671
                        ln.mu.Unlock()
17,947✔
2672
                }
2673
        }
2674
}
2675

2676
// updateLeafNodes will make sure to update the account smap for the subscription.
2677
// Will also forward to all leaf nodes as needed.
2678
func (acc *Account) updateLeafNodes(sub *subscription, delta int32) {
2,503,710✔
2679
        acc.updateLeafNodesEx(sub, delta, false)
2,503,710✔
2680
}
2,503,710✔
2681

2682
// This will make an update to our internal smap and determine if we should send out
2683
// an interest update to the remote side.
2684
// Lock should be held.
2685
func (c *client) updateSmap(sub *subscription, delta int32, isLDS bool) {
17,947✔
2686
        if c.leaf.smap == nil {
17,958✔
2687
                return
11✔
2688
        }
11✔
2689

2690
        // If we are solicited make sure this is a local client or a non-solicited leaf node
2691
        skind := sub.client.kind
17,936✔
2692
        updateClient := skind == CLIENT || skind == SYSTEM || skind == JETSTREAM || skind == ACCOUNT
17,936✔
2693
        if !isLDS && c.isSpokeLeafNode() && !(updateClient || (skind == LEAF && !sub.client.isSpokeLeafNode())) {
23,210✔
2694
                return
5,274✔
2695
        }
5,274✔
2696

2697
        // For additions, check if that sub has just been processed during initLeafNodeSmapAndSendSubs
2698
        if delta > 0 && c.leaf.tsub != nil {
18,693✔
2699
                if _, present := c.leaf.tsub[sub]; present {
6,032✔
2700
                        delete(c.leaf.tsub, sub)
1✔
2701
                        if len(c.leaf.tsub) == 0 {
1✔
2702
                                c.leaf.tsub = nil
×
2703
                                c.leaf.tsubt.Stop()
×
2704
                                c.leaf.tsubt = nil
×
2705
                        }
×
2706
                        return
1✔
2707
                }
2708
        }
2709

2710
        key := keyFromSub(sub)
12,661✔
2711
        n, ok := c.leaf.smap[key]
12,661✔
2712
        if delta < 0 && !ok {
13,427✔
2713
                return
766✔
2714
        }
766✔
2715

2716
        // We will update if its a queue, if count is zero (or negative), or we were 0 and are N > 0.
2717
        update := sub.queue != nil || (n <= 0 && n+delta > 0) || (n > 0 && n+delta <= 0)
11,895✔
2718
        n += delta
11,895✔
2719
        if n > 0 {
20,430✔
2720
                c.leaf.smap[key] = n
8,535✔
2721
        } else {
11,895✔
2722
                delete(c.leaf.smap, key)
3,360✔
2723
        }
3,360✔
2724
        if update {
19,311✔
2725
                c.sendLeafNodeSubUpdate(key, n)
7,416✔
2726
        }
7,416✔
2727
}
2728

2729
// Used to force add subjects to the subject map.
2730
func (c *client) forceAddToSmap(subj string) {
2✔
2731
        c.mu.Lock()
2✔
2732
        defer c.mu.Unlock()
2✔
2733

2✔
2734
        if c.leaf.smap == nil {
2✔
2735
                return
×
2736
        }
×
2737
        n := c.leaf.smap[subj]
2✔
2738
        if n != 0 {
2✔
2739
                return
×
2740
        }
×
2741
        // Place into the map since it was not there.
2742
        c.leaf.smap[subj] = 1
2✔
2743
        c.sendLeafNodeSubUpdate(subj, 1)
2✔
2744
}
2745

2746
// Used to force remove a subject from the subject map.
2747
func (c *client) forceRemoveFromSmap(subj string) {
×
2748
        c.mu.Lock()
×
2749
        defer c.mu.Unlock()
×
2750

×
2751
        if c.leaf.smap == nil {
×
2752
                return
×
2753
        }
×
2754
        n := c.leaf.smap[subj]
×
2755
        if n == 0 {
×
2756
                return
×
2757
        }
×
2758
        n--
×
2759
        if n == 0 {
×
2760
                // Remove is now zero
×
2761
                delete(c.leaf.smap, subj)
×
2762
                c.sendLeafNodeSubUpdate(subj, 0)
×
2763
        } else {
×
2764
                c.leaf.smap[subj] = n
×
2765
        }
×
2766
}
2767

2768
// Send the subscription interest change to the other side.
2769
// Lock should be held.
2770
func (c *client) sendLeafNodeSubUpdate(key string, n int32) {
7,418✔
2771
        // If we are a spoke, we need to check if we are allowed to send this subscription over to the hub.
7,418✔
2772
        if c.isSpokeLeafNode() {
9,434✔
2773
                checkPerms := true
2,016✔
2774
                if len(key) > 0 && (key[0] == '$' || key[0] == '_') {
3,213✔
2775
                        if strings.HasPrefix(key, leafNodeLoopDetectionSubjectPrefix) ||
1,197✔
2776
                                strings.HasPrefix(key, oldGWReplyPrefix) ||
1,197✔
2777
                                strings.HasPrefix(key, gwReplyPrefix) {
1,257✔
2778
                                checkPerms = false
60✔
2779
                        }
60✔
2780
                }
2781
                if checkPerms {
3,972✔
2782
                        var subject string
1,956✔
2783
                        if sep := strings.IndexByte(key, ' '); sep != -1 {
2,381✔
2784
                                subject = key[:sep]
425✔
2785
                        } else {
1,956✔
2786
                                subject = key
1,531✔
2787
                        }
1,531✔
2788
                        if !c.canSubscribe(subject) {
1,956✔
2789
                                return
×
2790
                        }
×
2791
                }
2792
        }
2793
        // If we are here we can send over to the other side.
2794
        _b := [64]byte{}
7,418✔
2795
        b := bytes.NewBuffer(_b[:0])
7,418✔
2796
        c.writeLeafSub(b, key, n)
7,418✔
2797
        c.enqueueProto(b.Bytes())
7,418✔
2798
}
2799

2800
// Helper function to build the key.
2801
func keyFromSub(sub *subscription) string {
40,925✔
2802
        var sb strings.Builder
40,925✔
2803
        sb.Grow(len(sub.subject) + len(sub.queue) + 1)
40,925✔
2804
        sb.Write(sub.subject)
40,925✔
2805
        if sub.queue != nil {
43,052✔
2806
                // Just make the key subject spc group, e.g. 'foo bar'
2,127✔
2807
                sb.WriteByte(' ')
2,127✔
2808
                sb.Write(sub.queue)
2,127✔
2809
        }
2,127✔
2810
        return sb.String()
40,925✔
2811
}
2812

2813
const (
2814
        keyRoutedSub         = "R"
2815
        keyRoutedSubByte     = 'R'
2816
        keyRoutedLeafSub     = "L"
2817
        keyRoutedLeafSubByte = 'L'
2818
)
2819

2820
// Helper function to build the key that prevents collisions between normal
2821
// routed subscriptions and routed subscriptions on behalf of a leafnode.
2822
// Keys will look like this:
2823
// "R foo"          -> plain routed sub on "foo"
2824
// "R foo bar"      -> queue routed sub on "foo", queue "bar"
2825
// "L foo bar"      -> plain routed leaf sub on "foo", leaf "bar"
2826
// "L foo bar baz"  -> queue routed sub on "foo", queue "bar", leaf "baz"
2827
func keyFromSubWithOrigin(sub *subscription) string {
718,271✔
2828
        var sb strings.Builder
718,271✔
2829
        sb.Grow(2 + len(sub.origin) + 1 + len(sub.subject) + 1 + len(sub.queue))
718,271✔
2830
        leaf := len(sub.origin) > 0
718,271✔
2831
        if leaf {
732,570✔
2832
                sb.WriteByte(keyRoutedLeafSubByte)
14,299✔
2833
        } else {
718,271✔
2834
                sb.WriteByte(keyRoutedSubByte)
703,972✔
2835
        }
703,972✔
2836
        sb.WriteByte(' ')
718,271✔
2837
        sb.Write(sub.subject)
718,271✔
2838
        if sub.queue != nil {
744,357✔
2839
                sb.WriteByte(' ')
26,086✔
2840
                sb.Write(sub.queue)
26,086✔
2841
        }
26,086✔
2842
        if leaf {
732,570✔
2843
                sb.WriteByte(' ')
14,299✔
2844
                sb.Write(sub.origin)
14,299✔
2845
        }
14,299✔
2846
        return sb.String()
718,271✔
2847
}
2848

2849
// Lock should be held.
2850
func (c *client) writeLeafSub(w *bytes.Buffer, key string, n int32) {
31,047✔
2851
        if key == _EMPTY_ {
31,047✔
2852
                return
×
2853
        }
×
2854
        if n > 0 {
58,734✔
2855
                w.WriteString("LS+ " + key)
27,687✔
2856
                // Check for queue semantics, if found write n.
27,687✔
2857
                if strings.Contains(key, " ") {
28,426✔
2858
                        w.WriteString(" ")
739✔
2859
                        var b [12]byte
739✔
2860
                        var i = len(b)
739✔
2861
                        for l := n; l > 0; l /= 10 {
1,478✔
2862
                                i--
739✔
2863
                                b[i] = digits[l%10]
739✔
2864
                        }
739✔
2865
                        w.Write(b[i:])
739✔
2866
                        if c.trace {
739✔
2867
                                arg := fmt.Sprintf("%s %d", key, n)
×
2868
                                c.traceOutOp("LS+", []byte(arg))
×
2869
                        }
×
2870
                } else if c.trace {
26,965✔
2871
                        c.traceOutOp("LS+", []byte(key))
17✔
2872
                }
17✔
2873
        } else {
3,360✔
2874
                w.WriteString("LS- " + key)
3,360✔
2875
                if c.trace {
3,360✔
2876
                        c.traceOutOp("LS-", []byte(key))
×
2877
                }
×
2878
        }
2879
        w.WriteString(CR_LF)
31,047✔
2880
}
2881

2882
// processLeafSub will process an inbound sub request for the remote leaf node.
2883
func (c *client) processLeafSub(argo []byte) (err error) {
27,432✔
2884
        // Indicate activity.
27,432✔
2885
        c.in.subs++
27,432✔
2886

27,432✔
2887
        srv := c.srv
27,432✔
2888
        if srv == nil {
27,432✔
2889
                return nil
×
2890
        }
×
2891

2892
        // Copy so we do not reference a potentially large buffer
2893
        arg := make([]byte, len(argo))
27,432✔
2894
        copy(arg, argo)
27,432✔
2895

27,432✔
2896
        args := splitArg(arg)
27,432✔
2897
        sub := &subscription{client: c}
27,432✔
2898

27,432✔
2899
        delta := int32(1)
27,432✔
2900
        switch len(args) {
27,432✔
2901
        case 1:
26,706✔
2902
                sub.queue = nil
26,706✔
2903
        case 3:
726✔
2904
                sub.queue = args[1]
726✔
2905
                sub.qw = int32(parseSize(args[2]))
726✔
2906
                // TODO: (ik) We should have a non empty queue name and a queue
726✔
2907
                // weight >= 1. For 2.11, we may want to return an error if that
726✔
2908
                // is not the case, but for now just overwrite `delta` if queue
726✔
2909
                // weight is greater than 1 (it is possible after a reconnect/
726✔
2910
                // server restart to receive a queue weight > 1 for a new sub).
726✔
2911
                if sub.qw > 1 {
944✔
2912
                        delta = sub.qw
218✔
2913
                }
218✔
2914
        default:
×
2915
                return fmt.Errorf("processLeafSub Parse Error: '%s'", arg)
×
2916
        }
2917
        sub.subject = args[0]
27,432✔
2918

27,432✔
2919
        c.mu.Lock()
27,432✔
2920
        if c.isClosed() {
27,450✔
2921
                c.mu.Unlock()
18✔
2922
                return nil
18✔
2923
        }
18✔
2924

2925
        acc := c.acc
27,414✔
2926
        // Guard against LS+ arriving before CONNECT has been processed, which
27,414✔
2927
        // can happen when compression is enabled.
27,414✔
2928
        if acc == nil {
27,414✔
2929
                c.mu.Unlock()
×
2930
                c.sendErr("Authorization Violation")
×
2931
                c.closeConnection(ProtocolViolation)
×
2932
                return nil
×
2933
        }
×
2934
        // Check if we have a loop.
2935
        ldsPrefix := bytes.HasPrefix(sub.subject, []byte(leafNodeLoopDetectionSubjectPrefix))
27,414✔
2936

27,414✔
2937
        if ldsPrefix && bytesToString(sub.subject) == acc.getLDSubject() {
27,418✔
2938
                c.mu.Unlock()
4✔
2939
                c.handleLeafNodeLoop(true)
4✔
2940
                return nil
4✔
2941
        }
4✔
2942

2943
        // Check permissions if applicable. (but exclude the $LDS, $GR and _GR_)
2944
        checkPerms := true
27,410✔
2945
        if sub.subject[0] == '$' || sub.subject[0] == '_' {
53,563✔
2946
                if ldsPrefix ||
26,153✔
2947
                        bytes.HasPrefix(sub.subject, []byte(oldGWReplyPrefix)) ||
26,153✔
2948
                        bytes.HasPrefix(sub.subject, []byte(gwReplyPrefix)) {
27,807✔
2949
                        checkPerms = false
1,654✔
2950
                }
1,654✔
2951
        }
2952

2953
        // If we are a hub check that we can publish to this subject.
2954
        if checkPerms {
53,166✔
2955
                subj := string(sub.subject)
25,756✔
2956
                if subjectIsLiteral(subj) && !c.pubAllowedFullCheck(subj, true, true) {
25,773✔
2957
                        c.mu.Unlock()
17✔
2958
                        c.leafSubPermViolation(sub.subject)
17✔
2959
                        c.Debugf(fmt.Sprintf("Permissions Violation for Subscription to %q", sub.subject))
17✔
2960
                        return nil
17✔
2961
                }
17✔
2962
        }
2963

2964
        // Check if we have a maximum on the number of subscriptions.
2965
        if c.subsAtLimit() {
27,401✔
2966
                c.mu.Unlock()
8✔
2967
                c.maxSubsExceeded()
8✔
2968
                return nil
8✔
2969
        }
8✔
2970

2971
        // If we have an origin cluster associated mark that in the sub.
2972
        if rc := c.remoteCluster(); rc != _EMPTY_ {
51,450✔
2973
                sub.origin = []byte(rc)
24,065✔
2974
        }
24,065✔
2975

2976
        // Like Routes, we store local subs by account and subject and optionally queue name.
2977
        // If we have a queue it will have a trailing weight which we do not want.
2978
        if sub.queue != nil {
28,109✔
2979
                sub.sid = arg[:len(arg)-len(args[2])-1]
724✔
2980
        } else {
27,385✔
2981
                sub.sid = arg
26,661✔
2982
        }
26,661✔
2983
        key := bytesToString(sub.sid)
27,385✔
2984
        osub := c.subs[key]
27,385✔
2985
        if osub == nil {
54,460✔
2986
                c.subs[key] = sub
27,075✔
2987
                // Now place into the account sl.
27,075✔
2988
                if err := acc.sl.Insert(sub); err != nil {
27,075✔
2989
                        delete(c.subs, key)
×
2990
                        c.mu.Unlock()
×
2991
                        c.Errorf("Could not insert subscription: %v", err)
×
2992
                        c.sendErr("Invalid Subscription")
×
2993
                        return nil
×
2994
                }
×
2995
        } else if sub.queue != nil {
619✔
2996
                // For a queue we need to update the weight.
309✔
2997
                delta = sub.qw - atomic.LoadInt32(&osub.qw)
309✔
2998
                atomic.StoreInt32(&osub.qw, sub.qw)
309✔
2999
                acc.sl.UpdateRemoteQSub(osub)
309✔
3000
        }
309✔
3001
        spoke := c.isSpokeLeafNode()
27,385✔
3002
        c.mu.Unlock()
27,385✔
3003

27,385✔
3004
        // Only add in shadow subs if a new sub or qsub.
27,385✔
3005
        if osub == nil {
54,460✔
3006
                if err := c.addShadowSubscriptions(acc, sub); err != nil {
27,075✔
3007
                        c.Errorf(err.Error())
×
3008
                }
×
3009
        }
3010

3011
        // If we are not solicited, treat leaf node subscriptions similar to a
3012
        // client subscription, meaning we forward them to routes, gateways and
3013
        // other leaf nodes as needed.
3014
        if !spoke {
37,173✔
3015
                // If we are routing add to the route map for the associated account.
9,788✔
3016
                srv.updateRouteSubscriptionMap(acc, sub, delta)
9,788✔
3017
                if srv.gateway.enabled {
10,687✔
3018
                        srv.gatewayUpdateSubInterest(acc.Name, sub, delta)
899✔
3019
                }
899✔
3020
        }
3021
        // Now check on leafnode updates for other leaf nodes. We understand solicited
3022
        // and non-solicited state in this call so we will do the right thing.
3023
        acc.updateLeafNodes(sub, delta)
27,385✔
3024

27,385✔
3025
        return nil
27,385✔
3026
}
3027

3028
// If the leafnode is a solicited, set the connect delay based on default
3029
// or private option (for tests). Sends the error to the other side, log and
3030
// close the connection.
3031
func (c *client) handleLeafNodeLoop(sendErr bool) {
12✔
3032
        accName, delay := c.setLeafConnectDelayIfSoliciting(leafNodeReconnectDelayAfterLoopDetected)
12✔
3033
        errTxt := fmt.Sprintf("Loop detected for leafnode account=%q. Delaying attempt to reconnect for %v", accName, delay)
12✔
3034
        if sendErr {
18✔
3035
                c.sendErr(errTxt)
6✔
3036
        }
6✔
3037

3038
        c.Errorf(errTxt)
12✔
3039
        // If we are here with "sendErr" false, it means that this is the server
12✔
3040
        // that received the error. The other side will have closed the connection,
12✔
3041
        // but does not hurt to close here too.
12✔
3042
        c.closeConnection(ProtocolViolation)
12✔
3043
}
3044

3045
// processLeafUnsub will process an inbound unsub request for the remote leaf node.
3046
func (c *client) processLeafUnsub(arg []byte) error {
3,011✔
3047
        // Indicate any activity, so pub and sub or unsubs.
3,011✔
3048
        c.in.subs++
3,011✔
3049

3,011✔
3050
        srv := c.srv
3,011✔
3051

3,011✔
3052
        c.mu.Lock()
3,011✔
3053
        if c.isClosed() {
3,050✔
3054
                c.mu.Unlock()
39✔
3055
                return nil
39✔
3056
        }
39✔
3057

3058
        acc := c.acc
2,972✔
3059
        // Guard against LS- arriving before CONNECT has been processed.
2,972✔
3060
        if acc == nil {
2,972✔
3061
                c.mu.Unlock()
×
3062
                c.sendErr("Authorization Violation")
×
3063
                c.closeConnection(ProtocolViolation)
×
3064
                return nil
×
3065
        }
×
3066

3067
        spoke := c.isSpokeLeafNode()
2,972✔
3068
        // We store local subs by account and subject and optionally queue name.
2,972✔
3069
        // LS- will have the arg exactly as the key.
2,972✔
3070
        sub, ok := c.subs[string(arg)]
2,972✔
3071
        if !ok {
2,972✔
3072
                // If not found, don't try to update routes/gws/leaf nodes.
×
3073
                c.mu.Unlock()
×
3074
                return nil
×
3075
        }
×
3076
        delta := int32(1)
2,972✔
3077
        if len(sub.queue) > 0 {
3,356✔
3078
                delta = sub.qw
384✔
3079
        }
384✔
3080
        c.mu.Unlock()
2,972✔
3081

2,972✔
3082
        c.unsubscribe(acc, sub, true, true)
2,972✔
3083
        if !spoke {
3,864✔
3084
                // If we are routing subtract from the route map for the associated account.
892✔
3085
                srv.updateRouteSubscriptionMap(acc, sub, -delta)
892✔
3086
                // Gateways
892✔
3087
                if srv.gateway.enabled {
1,069✔
3088
                        srv.gatewayUpdateSubInterest(acc.Name, sub, -delta)
177✔
3089
                }
177✔
3090
        }
3091
        // Now check on leafnode updates for other leaf nodes.
3092
        acc.updateLeafNodes(sub, -delta)
2,972✔
3093
        return nil
2,972✔
3094
}
3095

3096
func (c *client) processLeafHeaderMsgArgs(arg []byte) error {
223✔
3097
        // Unroll splitArgs to avoid runtime/heap issues
223✔
3098
        args := c.argsa[:0]
223✔
3099
        start := -1
223✔
3100
        for i, b := range arg {
12,331✔
3101
                switch b {
12,108✔
3102
                case ' ', '\t', '\r', '\n':
654✔
3103
                        if start >= 0 {
1,308✔
3104
                                args = append(args, arg[start:i])
654✔
3105
                                start = -1
654✔
3106
                        }
654✔
3107
                default:
11,454✔
3108
                        if start < 0 {
12,331✔
3109
                                start = i
877✔
3110
                        }
877✔
3111
                }
3112
        }
3113
        if start >= 0 {
446✔
3114
                args = append(args, arg[start:])
223✔
3115
        }
223✔
3116

3117
        c.pa.arg = arg
223✔
3118
        switch len(args) {
223✔
3119
        case 0, 1, 2:
×
3120
                return fmt.Errorf("processLeafHeaderMsgArgs Parse Error: '%s'", args)
×
3121
        case 3:
20✔
3122
                c.pa.reply = nil
20✔
3123
                c.pa.queues = nil
20✔
3124
                c.pa.hdb = args[1]
20✔
3125
                c.pa.hdr = parseSize(args[1])
20✔
3126
                c.pa.szb = args[2]
20✔
3127
                c.pa.size = parseSize(args[2])
20✔
3128
        case 4:
200✔
3129
                c.pa.reply = args[1]
200✔
3130
                c.pa.queues = nil
200✔
3131
                c.pa.hdb = args[2]
200✔
3132
                c.pa.hdr = parseSize(args[2])
200✔
3133
                c.pa.szb = args[3]
200✔
3134
                c.pa.size = parseSize(args[3])
200✔
3135
        default:
3✔
3136
                // args[1] is our reply indicator. Should be + or | normally.
3✔
3137
                if len(args[1]) != 1 {
3✔
3138
                        return fmt.Errorf("processLeafHeaderMsgArgs Bad or Missing Reply Indicator: '%s'", args[1])
×
3139
                }
×
3140
                switch args[1][0] {
3✔
3141
                case '+':
2✔
3142
                        c.pa.reply = args[2]
2✔
3143
                case '|':
1✔
3144
                        c.pa.reply = nil
1✔
3145
                default:
×
3146
                        return fmt.Errorf("processLeafHeaderMsgArgs Bad or Missing Reply Indicator: '%s'", args[1])
×
3147
                }
3148
                // Grab header size.
3149
                c.pa.hdb = args[len(args)-2]
3✔
3150
                c.pa.hdr = parseSize(c.pa.hdb)
3✔
3151

3✔
3152
                // Grab size.
3✔
3153
                c.pa.szb = args[len(args)-1]
3✔
3154
                c.pa.size = parseSize(c.pa.szb)
3✔
3155

3✔
3156
                // Grab queue names.
3✔
3157
                if c.pa.reply != nil {
5✔
3158
                        c.pa.queues = args[3 : len(args)-2]
2✔
3159
                } else {
3✔
3160
                        c.pa.queues = args[2 : len(args)-2]
1✔
3161
                }
1✔
3162
        }
3163
        if c.pa.hdr < 0 {
223✔
3164
                return fmt.Errorf("processLeafHeaderMsgArgs Bad or Missing Header Size: '%s'", arg)
×
3165
        }
×
3166
        if c.pa.size < 0 {
223✔
3167
                return fmt.Errorf("processLeafHeaderMsgArgs Bad or Missing Size: '%s'", args)
×
3168
        }
×
3169
        if c.pa.hdr > c.pa.size {
223✔
3170
                return fmt.Errorf("processLeafHeaderMsgArgs Header Size larger then TotalSize: '%s'", arg)
×
3171
        }
×
3172
        maxPayload := atomic.LoadInt32(&c.mpay)
223✔
3173
        if maxPayload != jwt.NoLimit && int64(c.pa.size) > int64(maxPayload) {
223✔
3174
                c.maxPayloadViolation(c.pa.size, maxPayload)
×
3175
                return ErrMaxPayload
×
3176
        }
×
3177

3178
        // Common ones processed after check for arg length
3179
        c.pa.subject = args[0]
223✔
3180

223✔
3181
        return nil
223✔
3182
}
3183

3184
func (c *client) processLeafMsgArgs(arg []byte) error {
58,361✔
3185
        // Unroll splitArgs to avoid runtime/heap issues
58,361✔
3186
        args := c.argsa[:0]
58,361✔
3187
        start := -1
58,361✔
3188
        for i, b := range arg {
2,063,282✔
3189
                switch b {
2,004,921✔
3190
                case ' ', '\t', '\r', '\n':
95,105✔
3191
                        if start >= 0 {
190,210✔
3192
                                args = append(args, arg[start:i])
95,105✔
3193
                                start = -1
95,105✔
3194
                        }
95,105✔
3195
                default:
1,909,816✔
3196
                        if start < 0 {
2,063,282✔
3197
                                start = i
153,466✔
3198
                        }
153,466✔
3199
                }
3200
        }
3201
        if start >= 0 {
116,722✔
3202
                args = append(args, arg[start:])
58,361✔
3203
        }
58,361✔
3204

3205
        c.pa.arg = arg
58,361✔
3206
        switch len(args) {
58,361✔
3207
        case 0, 1:
×
3208
                return fmt.Errorf("processLeafMsgArgs Parse Error: '%s'", args)
×
3209
        case 2:
36,976✔
3210
                c.pa.reply = nil
36,976✔
3211
                c.pa.queues = nil
36,976✔
3212
                c.pa.szb = args[1]
36,976✔
3213
                c.pa.size = parseSize(args[1])
36,976✔
3214
        case 3:
6,184✔
3215
                c.pa.reply = args[1]
6,184✔
3216
                c.pa.queues = nil
6,184✔
3217
                c.pa.szb = args[2]
6,184✔
3218
                c.pa.size = parseSize(args[2])
6,184✔
3219
        default:
15,201✔
3220
                // args[1] is our reply indicator. Should be + or | normally.
15,201✔
3221
                if len(args[1]) != 1 {
15,201✔
3222
                        return fmt.Errorf("processLeafMsgArgs Bad or Missing Reply Indicator: '%s'", args[1])
×
3223
                }
×
3224
                switch args[1][0] {
15,201✔
3225
                case '+':
158✔
3226
                        c.pa.reply = args[2]
158✔
3227
                case '|':
15,043✔
3228
                        c.pa.reply = nil
15,043✔
3229
                default:
×
3230
                        return fmt.Errorf("processLeafMsgArgs Bad or Missing Reply Indicator: '%s'", args[1])
×
3231
                }
3232
                // Grab size.
3233
                c.pa.szb = args[len(args)-1]
15,201✔
3234
                c.pa.size = parseSize(c.pa.szb)
15,201✔
3235

15,201✔
3236
                // Grab queue names.
15,201✔
3237
                if c.pa.reply != nil {
15,359✔
3238
                        c.pa.queues = args[3 : len(args)-1]
158✔
3239
                } else {
15,201✔
3240
                        c.pa.queues = args[2 : len(args)-1]
15,043✔
3241
                }
15,043✔
3242
        }
3243
        if c.pa.size < 0 {
58,361✔
3244
                return fmt.Errorf("processLeafMsgArgs Bad or Missing Size: '%s'", args)
×
3245
        }
×
3246
        maxPayload := atomic.LoadInt32(&c.mpay)
58,361✔
3247
        if maxPayload != jwt.NoLimit && int64(c.pa.size) > int64(maxPayload) {
58,361✔
3248
                c.maxPayloadViolation(c.pa.size, maxPayload)
×
3249
                return ErrMaxPayload
×
3250
        }
×
3251

3252
        // Common ones processed after check for arg length
3253
        c.pa.subject = args[0]
58,361✔
3254

58,361✔
3255
        return nil
58,361✔
3256
}
3257

3258
// processInboundLeafMsg is called to process an inbound msg from a leaf node.
3259
func (c *client) processInboundLeafMsg(msg []byte) {
57,264✔
3260
        // Update statistics
57,264✔
3261
        // The msg includes the CR_LF, so pull back out for accounting.
57,264✔
3262
        c.in.msgs++
57,264✔
3263
        c.in.bytes += int32(len(msg) - LEN_CR_LF)
57,264✔
3264

57,264✔
3265
        srv, acc, subject := c.srv, c.acc, string(c.pa.subject)
57,264✔
3266

57,264✔
3267
        // Mostly under testing scenarios.
57,264✔
3268
        if srv == nil || acc == nil {
57,264✔
3269
                return
×
3270
        }
×
3271

3272
        // Check that leaf messages respect the subject permissions.
3273
        if c.perms != nil && !c.leafMsgAllowed() {
57,264✔
3274
                c.leafPubPermViolation(c.pa.subject)
×
3275
                return
×
3276
        }
×
3277

3278
        // Match the subscriptions. We will use our own L1 map if
3279
        // it's still valid, avoiding contention on the shared sublist.
3280
        var r *SublistResult
57,264✔
3281
        var ok bool
57,264✔
3282

57,264✔
3283
        genid := atomic.LoadUint64(&c.acc.sl.genid)
57,264✔
3284
        if genid == c.in.genid && c.in.results != nil {
112,577✔
3285
                r, ok = c.in.results[subject]
55,313✔
3286
        } else {
57,264✔
3287
                // Reset our L1 completely.
1,951✔
3288
                c.in.results = make(map[string]*SublistResult)
1,951✔
3289
                c.in.genid = genid
1,951✔
3290
        }
1,951✔
3291

3292
        // Go back to the sublist data structure.
3293
        if !ok {
92,326✔
3294
                r = c.acc.sl.Match(subject)
35,062✔
3295
                // Prune the results cache. Keeps us from unbounded growth. Random delete.
35,062✔
3296
                if len(c.in.results) >= maxResultCacheSize {
35,973✔
3297
                        n := 0
911✔
3298
                        for subj := range c.in.results {
30,974✔
3299
                                delete(c.in.results, subj)
30,063✔
3300
                                if n++; n > pruneSize {
30,974✔
3301
                                        break
911✔
3302
                                }
3303
                        }
3304
                }
3305
                // Then add the new cache entry.
3306
                c.in.results[subject] = r
35,062✔
3307
        }
3308

3309
        // Collect queue names if needed.
3310
        var qnames [][]byte
57,264✔
3311

57,264✔
3312
        // Check for no interest, short circuit if so.
57,264✔
3313
        // This is the fanout scale.
57,264✔
3314
        if len(r.psubs)+len(r.qsubs) > 0 {
114,233✔
3315
                flag := pmrNoFlag
56,969✔
3316
                // If we have queue subs in this cluster, then if we run in gateway
56,969✔
3317
                // mode and the remote gateways have queue subs, then we need to
56,969✔
3318
                // collect the queue groups this message was sent to so that we
56,969✔
3319
                // exclude them when sending to gateways.
56,969✔
3320
                if len(r.qsubs) > 0 && c.srv.gateway.enabled &&
56,969✔
3321
                        atomic.LoadInt64(&c.srv.gateway.totalQSubs) > 0 {
62,232✔
3322
                        flag |= pmrCollectQueueNames
5,263✔
3323
                }
5,263✔
3324
                // If this is a mapped subject that means the mapped interest
3325
                // is what got us here, but this might not have a queue designation
3326
                // If that is the case, make sure we ignore to process local queue subscribers.
3327
                if len(c.pa.mapped) > 0 && len(c.pa.queues) == 0 {
57,213✔
3328
                        flag |= pmrIgnoreEmptyQueueFilter
244✔
3329
                }
244✔
3330
                _, qnames = c.processMsgResults(acc, r, msg, nil, c.pa.subject, c.pa.reply, flag)
56,969✔
3331
        }
3332

3333
        // Now deal with gateways
3334
        if c.srv.gateway.enabled {
63,303✔
3335
                c.sendMsgToGateways(acc, msg, c.pa.subject, c.pa.reply, qnames, true)
6,039✔
3336
        }
6,039✔
3337
}
3338

3339
// Checks whether the inbound leaf message is allowed by the
3340
// connection's permissions. On the hub side this enforces what
3341
// the remote leaf may publish. On the spoke side this enforces
3342
// import restrictions such as deny_imports.
3343
func (c *client) leafMsgAllowed() bool {
53,807✔
3344
        wireSubject := c.pa.subject
53,807✔
3345
        if len(c.pa.mapped) > 0 {
54,051✔
3346
                // Mappings rewrite c.pa.subject to the internal
244✔
3347
                // destination. For leaf ACLs, need to check
244✔
3348
                // the original wire subject from the remote side.
244✔
3349
                wireSubject = c.pa.mapped
244✔
3350
        }
244✔
3351
        // Strip any gateway routing prefix for the permission check.
3352
        subjectToCheck, isGW := getGWRoutedSubjectOrSelf(wireSubject)
53,807✔
3353

53,807✔
3354
        // Service-import replies (_R_), JS ack subjects ($JS.ACK.)
53,807✔
3355
        // are internal routing subjects forwarded via LS+ without
53,807✔
3356
        // permission checks.
53,807✔
3357
        if isServiceReply(subjectToCheck) || isJSAckSubject(subjectToCheck) {
53,837✔
3358
                return true
30✔
3359
        }
30✔
3360

3361
        c.mu.RLock()
53,777✔
3362
        if c.isSpokeLeafNode() {
81,047✔
3363
                // Gateway routed replies are forwarded without
27,270✔
3364
                // permission checks.
27,270✔
3365
                if isGW || c.leafReceiveAllowed(subjectToCheck) {
54,540✔
3366
                        c.mu.RUnlock()
27,270✔
3367
                        return true
27,270✔
3368
                }
27,270✔
3369
        } else if c.leafSendAllowed(subjectToCheck) {
53,014✔
3370
                c.mu.RUnlock()
26,507✔
3371
                return true
26,507✔
3372
        }
26,507✔
3373

3374
        // If allow_responses is not configured, or there is no tracked reply for
3375
        // this subject, the answer is "denied" and we can return it while still
3376
        // holding only the read lock.
3377
        replySubject := bytesToString(wireSubject)
×
3378
        if c.perms == nil || c.perms.resp == nil || c.replies[replySubject] == nil {
×
3379
                c.mu.RUnlock()
×
3380
                return false
×
3381
        }
×
3382
        c.mu.RUnlock()
×
3383

×
3384
        // Check tracked reply permissions (allow_responses).
×
3385
        // Use the pre-strip subject since deliverMsg tracks
×
3386
        // replies under the original form, which includes
×
3387
        // the GW routing prefix for routed requests.
×
3388
        c.mu.Lock()
×
3389
        defer c.mu.Unlock()
×
3390
        return c.responseAllowed(replySubject)
×
3391
}
3392

3393
// Returns true if the leaf side ACLs allow importing this subject,
3394
// based on the permissions received over INFO and any local deny_imports.
3395
// At least a read lock must be held.
3396
func (c *client) leafReceiveAllowed(subject []byte) bool {
27,270✔
3397
        return c.canSubscribeInternal(bytesToString(subject))
27,270✔
3398
}
27,270✔
3399

3400
// Returns true if the hub side ACLs allow the remote leaf to send
3401
// this subject.
3402
// At least a read lock must be held.
3403
func (c *client) leafSendAllowed(bsubject []byte) bool {
26,507✔
3404
        // Use the original export ACL captured for this accepted leaf.
26,507✔
3405
        // The live perms also contain additional JetStream denies used by
26,507✔
3406
        // the normal forwarding path, and applying them here would reject
26,507✔
3407
        // legitimate inbound JS API requests.
26,507✔
3408
        subject := bytesToString(bsubject)
26,507✔
3409
        perms := c.opts.Export
26,507✔
3410
        if perms == nil || (perms.Allow == nil && perms.Deny == nil) {
53,011✔
3411
                return true
26,504✔
3412
        }
26,504✔
3413

3414
        allowed := true
3✔
3415
        if perms.Allow != nil && !strings.HasPrefix(subject, mqttPrefix) {
5✔
3416
                allowed = false
2✔
3417
                for _, allowSubj := range perms.Allow {
5✔
3418
                        if matchLiteral(subject, allowSubj) {
5✔
3419
                                allowed = true
2✔
3420
                                break
2✔
3421
                        }
3422
                }
3423
        }
3424

3425
        if allowed && len(perms.Deny) > 0 {
4✔
3426
                for _, denySubj := range perms.Deny {
2✔
3427
                        if matchLiteral(subject, denySubj) {
1✔
3428
                                allowed = false
×
3429
                                break
×
3430
                        }
3431
                }
3432
        }
3433
        return allowed
3✔
3434
}
3435

3436
// Handles a subscription permission violation.
3437
// See leafPermViolation() for details.
3438
func (c *client) leafSubPermViolation(subj []byte) {
17✔
3439
        c.leafPermViolation(false, subj)
17✔
3440
}
17✔
3441

3442
// Handles a publish permission violation.
3443
// See leafPermViolation() for details.
3444
func (c *client) leafPubPermViolation(subj []byte) {
×
3445
        c.leafPermViolation(true, subj)
×
3446
}
×
3447

3448
// Common function to process publish or subscribe leafnode permission violation.
3449
// Sends the permission violation error to the remote, logs it and closes the connection.
3450
// If this is from a server soliciting, the reconnection will be delayed.
3451
func (c *client) leafPermViolation(pub bool, subj []byte) {
17✔
3452
        if c.isSpokeLeafNode() {
34✔
3453
                // For spokes these are no-ops since the hub server told us our permissions.
17✔
3454
                // We just need to not send these over to the other side since we will get cutoff.
17✔
3455
                return
17✔
3456
        }
17✔
3457
        // FIXME(dlc) ?
3458
        c.setLeafConnectDelayIfSoliciting(leafNodeReconnectAfterPermViolation)
×
3459
        var action string
×
3460
        if pub {
×
3461
                c.sendErr(fmt.Sprintf("Permissions Violation for Publish to %q", subj))
×
3462
                action = "Publish"
×
3463
        } else {
×
3464
                c.sendErr(fmt.Sprintf("Permissions Violation for Subscription to %q", subj))
×
3465
                action = "Subscription"
×
3466
        }
×
3467
        c.Errorf("%s Violation on %q - Check other side configuration", action, subj)
×
3468
        // TODO: add a new close reason that is more appropriate?
×
3469
        c.closeConnection(ProtocolViolation)
×
3470
}
3471

3472
// Invoked from generic processErr() for LEAF connections.
3473
func (c *client) leafProcessErr(errStr string) {
46✔
3474
        // Check if we got a cluster name collision.
46✔
3475
        if strings.Contains(errStr, ErrLeafNodeHasSameClusterName.Error()) {
49✔
3476
                _, delay := c.setLeafConnectDelayIfSoliciting(leafNodeReconnectDelayAfterClusterNameSame)
3✔
3477
                c.Errorf("Leafnode connection dropped with same cluster name error. Delaying attempt to reconnect for %v", delay)
3✔
3478
                return
3✔
3479
        }
3✔
3480
        if strings.Contains(errStr, ErrLeafNodeMinVersionRejected.Error()) {
44✔
3481
                _, delay := c.setLeafConnectDelayIfSoliciting(leafNodeMinVersionReconnectDelay)
1✔
3482
                c.Errorf("Leafnode connection dropped due to minimum version requirement. Delaying attempt to reconnect for %v", delay)
1✔
3483
                return
1✔
3484
        }
1✔
3485

3486
        // We will look for Loop detected error coming from the other side.
3487
        // If we solicit, set the connect delay.
3488
        if !strings.Contains(errStr, "Loop detected") {
78✔
3489
                return
36✔
3490
        }
36✔
3491
        c.handleLeafNodeLoop(false)
6✔
3492
}
3493

3494
// If this leaf connection solicits, sets the connect delay to the given value,
3495
// or the one from the server option's LeafNode.connDelay if one is set (for tests).
3496
// Returns the connection's account name and delay.
3497
func (c *client) setLeafConnectDelayIfSoliciting(delay time.Duration) (string, time.Duration) {
16✔
3498
        c.mu.Lock()
16✔
3499
        if c.isSolicitedLeafNode() {
26✔
3500
                if s := c.srv; s != nil {
20✔
3501
                        if srvdelay := s.getOpts().LeafNode.connDelay; srvdelay != 0 {
14✔
3502
                                delay = srvdelay
4✔
3503
                        }
4✔
3504
                }
3505
                c.leaf.remote.setConnectDelay(delay)
10✔
3506
        }
3507
        var accName string
16✔
3508
        if c.acc != nil {
32✔
3509
                accName = c.acc.Name
16✔
3510
        }
16✔
3511
        c.mu.Unlock()
16✔
3512
        return accName, delay
16✔
3513
}
3514

3515
// For the given remote Leafnode configuration, this function returns
3516
// if TLS is required, and if so, will return a clone of the TLS Config
3517
// (since some fields will be changed during handshake), the TLS server
3518
// name that is remembered, and the TLS timeout.
3519
func (c *client) leafNodeGetTLSConfigForSolicit(remote *leafNodeCfg) (bool, *tls.Config, string, float64) {
1,536✔
3520
        var (
1,536✔
3521
                tlsConfig  *tls.Config
1,536✔
3522
                tlsName    string
1,536✔
3523
                tlsTimeout float64
1,536✔
3524
        )
1,536✔
3525

1,536✔
3526
        remote.RLock()
1,536✔
3527
        defer remote.RUnlock()
1,536✔
3528

1,536✔
3529
        tlsRequired := remote.TLS || remote.TLSConfig != nil
1,536✔
3530
        if tlsRequired {
1,610✔
3531
                if remote.TLSConfig != nil {
125✔
3532
                        tlsConfig = remote.TLSConfig.Clone()
51✔
3533
                } else {
74✔
3534
                        tlsConfig = &tls.Config{MinVersion: tls.VersionTLS12}
23✔
3535
                }
23✔
3536
                tlsName = remote.tlsName
74✔
3537
                tlsTimeout = remote.TLSTimeout
74✔
3538
                if tlsTimeout == 0 {
114✔
3539
                        tlsTimeout = float64(TLS_TIMEOUT / time.Second)
40✔
3540
                }
40✔
3541
        }
3542

3543
        return tlsRequired, tlsConfig, tlsName, tlsTimeout
1,536✔
3544
}
3545

3546
// Initiates the LeafNode Websocket connection by:
3547
// - doing the TLS handshake if needed
3548
// - sending the HTTP request
3549
// - waiting for the HTTP response
3550
//
3551
// Since some bufio reader is used to consume the HTTP response, this function
3552
// returns the slice of buffered bytes (if any) so that the readLoop that will
3553
// be started after that consume those first before reading from the socket.
3554
// The boolean
3555
//
3556
// Lock held on entry.
3557
func (c *client) leafNodeSolicitWSConnection(opts *Options, rURL *url.URL, remote *leafNodeCfg) ([]byte, ClosedState, error) {
46✔
3558
        remote.RLock()
46✔
3559
        compress := remote.Websocket.Compression
46✔
3560
        // By default the server will mask outbound frames, but it can be disabled with this option.
46✔
3561
        noMasking := remote.Websocket.NoMasking
46✔
3562
        infoTimeout := remote.FirstInfoTimeout
46✔
3563
        remote.RUnlock()
46✔
3564
        // Will do the client-side TLS handshake if needed.
46✔
3565
        tlsRequired, err := c.leafClientHandshakeIfNeeded(remote, opts)
46✔
3566
        if err != nil {
50✔
3567
                // 0 will indicate that the connection was already closed
4✔
3568
                return nil, 0, err
4✔
3569
        }
4✔
3570

3571
        // For http request, we need the passed URL to contain either http or https scheme.
3572
        scheme := "http"
42✔
3573
        if tlsRequired {
50✔
3574
                scheme = "https"
8✔
3575
        }
8✔
3576
        // We will use the `/leafnode` path to tell the accepting WS server that it should
3577
        // create a LEAF connection, not a CLIENT.
3578
        // In case we use the user's URL path in the future, make sure we append the user's
3579
        // path to our `/leafnode` path.
3580
        lpath := leafNodeWSPath
42✔
3581
        if curPath := rURL.EscapedPath(); curPath != _EMPTY_ {
63✔
3582
                if curPath[0] == '/' {
42✔
3583
                        curPath = curPath[1:]
21✔
3584
                }
21✔
3585
                lpath = path.Join(curPath, lpath)
21✔
3586
        } else {
21✔
3587
                lpath = lpath[1:]
21✔
3588
        }
21✔
3589
        ustr := fmt.Sprintf("%s://%s/%s", scheme, rURL.Host, lpath)
42✔
3590
        u, _ := url.Parse(ustr)
42✔
3591
        req := &http.Request{
42✔
3592
                Method:     "GET",
42✔
3593
                URL:        u,
42✔
3594
                Proto:      "HTTP/1.1",
42✔
3595
                ProtoMajor: 1,
42✔
3596
                ProtoMinor: 1,
42✔
3597
                Header:     make(http.Header),
42✔
3598
                Host:       u.Host,
42✔
3599
        }
42✔
3600
        wsKey, err := wsMakeChallengeKey()
42✔
3601
        if err != nil {
42✔
3602
                return nil, WriteError, err
×
3603
        }
×
3604

3605
        req.Header["Upgrade"] = []string{"websocket"}
42✔
3606
        req.Header["Connection"] = []string{"Upgrade"}
42✔
3607
        req.Header["Sec-WebSocket-Key"] = []string{wsKey}
42✔
3608
        req.Header["Sec-WebSocket-Version"] = []string{"13"}
42✔
3609
        if compress {
50✔
3610
                req.Header.Add("Sec-WebSocket-Extensions", wsPMCReqHeaderValue)
8✔
3611
        }
8✔
3612
        if noMasking {
51✔
3613
                req.Header.Add(wsNoMaskingHeader, wsNoMaskingValue)
9✔
3614
        }
9✔
3615
        c.nc.SetDeadline(time.Now().Add(infoTimeout))
42✔
3616
        if err := req.Write(c.nc); err != nil {
42✔
3617
                return nil, WriteError, err
×
3618
        }
×
3619

3620
        var resp *http.Response
42✔
3621

42✔
3622
        br := bufio.NewReaderSize(c.nc, MAX_CONTROL_LINE_SIZE)
42✔
3623
        resp, err = http.ReadResponse(br, req)
42✔
3624
        if err == nil &&
42✔
3625
                (resp.StatusCode != 101 ||
42✔
3626
                        !strings.EqualFold(resp.Header.Get("Upgrade"), "websocket") ||
42✔
3627
                        !strings.EqualFold(resp.Header.Get("Connection"), "upgrade") ||
42✔
3628
                        resp.Header.Get("Sec-Websocket-Accept") != wsAcceptKey(wsKey)) {
43✔
3629

1✔
3630
                err = fmt.Errorf("invalid websocket connection")
1✔
3631
        }
1✔
3632
        // Check compression extension...
3633
        if err == nil && c.ws.compress {
50✔
3634
                // Check that not only permessage-deflate extension is present, but that
8✔
3635
                // we also have server and client no context take over.
8✔
3636
                srvCompress, noCtxTakeover := wsPMCExtensionSupport(resp.Header, false)
8✔
3637

8✔
3638
                // If server does not support compression, then simply disable it in our side.
8✔
3639
                if !srvCompress {
12✔
3640
                        c.ws.compress = false
4✔
3641
                } else if !noCtxTakeover {
8✔
3642
                        err = fmt.Errorf("compression negotiation error")
×
3643
                }
×
3644
        }
3645
        // Same for no masking...
3646
        if err == nil && noMasking {
51✔
3647
                // Check if server accepts no masking
9✔
3648
                if resp.Header.Get(wsNoMaskingHeader) != wsNoMaskingValue {
10✔
3649
                        // Nope, need to mask our writes as any client would do.
1✔
3650
                        c.ws.maskwrite = true
1✔
3651
                }
1✔
3652
        }
3653
        if resp != nil {
69✔
3654
                resp.Body.Close()
27✔
3655
        }
27✔
3656
        if err != nil {
58✔
3657
                return nil, ReadError, err
16✔
3658
        }
16✔
3659
        c.Debugf("Leafnode compression=%v masking=%v", c.ws.compress, c.ws.maskwrite)
26✔
3660

26✔
3661
        var preBuf []byte
26✔
3662
        // We have to slurp whatever is in the bufio reader and pass that to the readloop.
26✔
3663
        if n := br.Buffered(); n != 0 {
26✔
3664
                preBuf, _ = br.Peek(n)
×
3665
        }
×
3666
        return preBuf, 0, nil
26✔
3667
}
3668

3669
const connectProcessTimeout = 2 * time.Second
3670

3671
// This is invoked for remote LEAF remote connections after processing the INFO
3672
// protocol.
3673
func (s *Server) leafNodeResumeConnectProcess(c *client) {
548✔
3674
        clusterName := s.ClusterName()
548✔
3675

548✔
3676
        c.mu.Lock()
548✔
3677
        if c.isClosed() {
548✔
3678
                c.mu.Unlock()
×
3679
                return
×
3680
        }
×
3681
        if err := c.sendLeafConnect(clusterName, c.headers); err != nil {
550✔
3682
                c.mu.Unlock()
2✔
3683
                c.closeConnection(WriteError)
2✔
3684
                return
2✔
3685
        }
2✔
3686

3687
        // Spin up the write loop.
3688
        s.startGoRoutine(func() { c.writeLoop() })
1,092✔
3689

3690
        // timeout leafNodeFinishConnectProcess
3691
        c.ping.tmr = time.AfterFunc(connectProcessTimeout, func() {
546✔
3692
                c.mu.Lock()
×
3693
                // check if leafNodeFinishConnectProcess was called and prevent later leafNodeFinishConnectProcess
×
3694
                if !c.flags.setIfNotSet(connectProcessFinished) {
×
3695
                        c.mu.Unlock()
×
3696
                        return
×
3697
                }
×
3698
                clearTimer(&c.ping.tmr)
×
3699
                closed := c.isClosed()
×
3700
                c.mu.Unlock()
×
3701
                if !closed {
×
3702
                        c.sendErrAndDebug("Stale Leaf Node Connection - Closing")
×
3703
                        c.closeConnection(StaleConnection)
×
3704
                }
×
3705
        })
3706
        c.mu.Unlock()
546✔
3707
        c.Debugf("Remote leafnode connect msg sent")
546✔
3708
}
3709

3710
// This is invoked for remote LEAF connections after processing the INFO
3711
// protocol and leafNodeResumeConnectProcess.
3712
// This will send LS+ the CONNECT protocol and register the leaf node.
3713
func (s *Server) leafNodeFinishConnectProcess(c *client) {
511✔
3714
        c.mu.Lock()
511✔
3715
        if !c.flags.setIfNotSet(connectProcessFinished) {
511✔
3716
                c.mu.Unlock()
×
3717
                return
×
3718
        }
×
3719
        if c.isClosed() {
511✔
3720
                c.mu.Unlock()
×
3721
                s.removeLeafNodeConnection(c)
×
3722
                return
×
3723
        }
×
3724
        remote := c.leaf.remote
511✔
3725
        if remote == nil || c.acc == nil {
511✔
3726
                c.mu.Unlock()
×
3727
                c.sendErr("Authorization Violation")
×
3728
                c.closeConnection(ProtocolViolation)
×
3729
                return
×
3730
        }
×
3731
        // Check if we will need to send the system connect event.
3732
        remote.RLock()
511✔
3733
        sendSysConnectEvent := remote.Hub
511✔
3734
        remote.RUnlock()
511✔
3735

511✔
3736
        // Capture account before releasing lock
511✔
3737
        acc := c.acc
511✔
3738
        // cancel connectProcessTimeout
511✔
3739
        clearTimer(&c.ping.tmr)
511✔
3740
        c.mu.Unlock()
511✔
3741

511✔
3742
        // Make sure we register with the account here.
511✔
3743
        if err := c.registerWithAccount(acc); err != nil {
513✔
3744
                if err == ErrTooManyAccountConnections {
2✔
3745
                        c.maxAccountConnExceeded()
×
3746
                        return
×
3747
                } else if err == ErrLeafNodeLoop {
4✔
3748
                        c.handleLeafNodeLoop(true)
2✔
3749
                        return
2✔
3750
                }
2✔
3751
                c.Errorf("Registering leaf with account %s resulted in error: %v", acc.Name, err)
×
3752
                c.closeConnection(ProtocolViolation)
×
3753
                return
×
3754
        }
3755
        if !s.addLeafNodeConnection(c, _EMPTY_, _EMPTY_, false) {
509✔
3756
                // Was not added, could be because the remote configuration has been removed.
×
3757
                c.closeConnection(ClientClosed)
×
3758
                return
×
3759
        }
×
3760
        s.initLeafNodeSmapAndSendSubs(c)
509✔
3761
        if sendSysConnectEvent {
521✔
3762
                s.sendLeafNodeConnect(acc)
12✔
3763
        }
12✔
3764
        s.accountConnectEvent(c)
509✔
3765

509✔
3766
        // The above functions are not running under the client lock, so it is
509✔
3767
        // possible that between the time we have started the read/write loops
509✔
3768
        // and now, that the connection was closed. This would leave the closed
509✔
3769
        // LN connection possibly registered with the account and/or the server's
509✔
3770
        // leafs map. So check if connection is closed, and if so, manually cleanup.
509✔
3771
        c.mu.Lock()
509✔
3772
        closed := c.isClosed()
509✔
3773
        if !closed {
1,017✔
3774
                c.setFirstPingTimer()
508✔
3775
        }
508✔
3776
        c.mu.Unlock()
509✔
3777
        if closed {
510✔
3778
                s.removeLeafNodeConnection(c)
1✔
3779
                if prev := acc.removeClient(c); prev == 1 {
2✔
3780
                        s.decActiveAccounts()
1✔
3781
                }
1✔
3782
        }
3783
}
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc