-
Notifications
You must be signed in to change notification settings - Fork 3
/
Copy pathmain.go
122 lines (103 loc) · 3.32 KB
/
main.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
// Copyright 2016-2018 The NATS Authors
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package main
import (
"flag"
"fmt"
"log"
"os"
"strconv"
"time"
nats "github.com/nats-io/nats.go"
"github.com/nats-io/stan.go"
)
var usageStr = `
Usage: stan-pub [options] <subject> <message>
Options:
-s, --server <url> NATS Streaming server URL(s)
-c, --cluster <cluster name> NATS Streaming cluster name
-cr, --creds <credentials> NATS 2.0 Credentials
-d, --delay <delay in milliseconds> Delay in publishing message in milliseonds.
`
// NOTE: Use tls scheme for TLS, e.g. stan-pub -s tls://demo.nats.io:4443 foo hello
func usage() {
fmt.Printf("%s\n", usageStr)
os.Exit(0)
}
func main() {
var (
clusterID string
clientID string
URL string
userCreds string
delay int
limit int
)
flag.StringVar(&URL, "s", stan.DefaultNatsURL, "The nats server URLs (separated by comma)")
flag.StringVar(&URL, "server", stan.DefaultNatsURL, "The nats server URLs (separated by comma)")
flag.StringVar(&clusterID, "c", "local-stan", "The NATS Streaming cluster ID")
flag.StringVar(&clusterID, "cluster", "local-stan", "The NATS Streaming cluster ID")
flag.StringVar(&userCreds, "cr", "", "Credentials File")
flag.StringVar(&userCreds, "creds", "", "Credentials File")
flag.IntVar(&delay, "d", 1000, "Delay in seconds between publishing message")
flag.IntVar(&delay, "delay", 1000, "Delay in seconds between publishing message")
flag.IntVar(&limit, "limit", -1, "Limit the number of messages to publish. Once this is reached, the publisher will no longer publish messages. -1 means no limit.")
log.SetFlags(0)
flag.Usage = usage
flag.Parse()
args := flag.Args()
if len(args) < 1 {
usage()
}
// Connect Options.
opts := []nats.Option{nats.Name("Gonuts Publisher")}
// Use UserCredentials
if userCreds != "" {
opts = append(opts, nats.UserCredentials(userCreds))
}
// Connect to NATS
nc, err := nats.Connect(URL, opts...)
if err != nil {
log.Fatal(err)
}
defer nc.Close()
clientID = strconv.FormatInt(time.Now().UnixNano(), 10)
sc, err := stan.Connect(clusterID, clientID, stan.NatsConn(nc))
if err != nil {
log.Fatalf("Can't connect: %v.\nMake sure a NATS Streaming Server is running at: %s", err, URL)
}
defer sc.Close()
subj := args[0]
msgPublishedCount := 0
var publishMsg = true
for {
if limit > 0 {
if msgPublishedCount >= limit {
publishMsg = false
}
}
if publishMsg {
t := time.Now()
msg := []byte("Message is : " + t.String())
err = sc.Publish(subj, msg)
if err != nil {
log.Fatalf("Error during publish: %v\n", err)
}
log.Printf("Published [%s] : '%s'\n", subj, msg)
time.Sleep(time.Duration(delay) * time.Millisecond)
msgPublishedCount++
} else {
log.Println("Reached publishing limit, stop publishing.")
}
}
}