-
Notifications
You must be signed in to change notification settings - Fork 0
/
publish.go
101 lines (89 loc) · 2.6 KB
/
publish.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
// Copyright 2024 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 (
"encoding/json"
"log"
"strconv"
"time"
)
// Message options
type messageOpts struct {
qos int
retain bool
size int
topic string
}
type publisher struct {
messageOpts
mps int
messages int
topics int
clientID string
}
func (p *publisher) publish(msgCh chan *Stat, errorCh chan error, timestamp bool) {
cl, _, cleanup, err := connect(p.clientID, CleanSession)
if err != nil {
log.Fatal(err)
}
defer cleanup()
opts := cl.OptionsReader()
start := time.Now()
var elapsed time.Duration
bc := 0
iTopic := 0
for n := 0; n < p.messages; n++ {
now := time.Now()
if n > 0 && p.mps > 0 {
next := start.Add(time.Duration(n) * time.Second / time.Duration(p.mps))
time.Sleep(next.Sub(now))
}
// payload always starts with JSON containing timestamp, etc. The JSON
// is always terminated with a '-', which can not be part of the random
// fill. payload is then filled to the requested size with random data.
payload := randomPayload(p.size)
if timestamp {
structuredPayload, _ := json.Marshal(PubValue{
Seq: n,
Timestamp: time.Now().UnixNano(),
})
structuredPayload = append(structuredPayload, '\n')
if len(structuredPayload) > len(payload) {
payload = structuredPayload
} else {
copy(payload, structuredPayload)
}
}
currTopic := p.topic
if p.topics > 0 {
currTopic = p.topic + "/" + strconv.Itoa(iTopic)
iTopic = (iTopic + 1) % p.topics
}
startPublish := time.Now()
if token := cl.Publish(currTopic, byte(p.qos), p.retain, payload); token.Wait() && token.Error() != nil {
errorCh <- token.Error()
return
}
elapsedPublish := time.Since(startPublish)
elapsed += elapsedPublish
logOp(opts.ClientID(), "PUB <-", elapsedPublish, "Published: %d bytes to %q, qos:%v, retain:%v", len(payload), currTopic, p.qos, p.retain)
bc += mqttPublishLen(currTopic, byte(p.qos), p.retain, payload)
}
if msgCh != nil {
msgCh <- &Stat{
Ops: p.messages,
NS: map[string]time.Duration{"pub": elapsed},
Bytes: int64(bc),
}
}
}