forked from VolantMQ/volantmq
-
Notifications
You must be signed in to change notification settings - Fork 0
/
clients.go
112 lines (95 loc) · 3.22 KB
/
clients.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
package systree
import (
"encoding/json"
"sync/atomic"
"time"
"github.com/VolantMQ/vlapi/mqttp"
"github.com/VolantMQ/volantmq/types"
)
// ClientConnectStatus is argument to client connected state
type ClientConnectStatus struct {
Address string
Username string
Timestamp string
ReceiveMaximum uint32
MaximumPacketSize uint32
KeepAlive uint16
GeneratedID bool
CleanSession bool
Durable bool
SessionPresent bool
PreserveOrder bool
MaximumQoS mqttp.QosType
Protocol mqttp.ProtocolVersion
ConnAckCode mqttp.ReasonCode
}
type clientDisconnectStatus struct {
Reason string
Timestamp string
}
type clients struct {
stat
topicsManager types.TopicMessenger
topic string
}
func newClients(topicPrefix string, retained *[]types.RetainObject) clients {
c := clients{
stat: newStat(topicPrefix+"/stats/clients", retained),
topic: topicPrefix + "/clients/",
}
return c
}
// Connected add to statistic new client
func (t *clients) Connected(id string, status *ClientConnectStatus) {
newVal := atomic.AddUint64(&t.curr.val, 1)
if atomic.LoadUint64(&t.max.val) < newVal {
atomic.StoreUint64(&t.max.val, newVal)
}
// notify client connected
nm, _ := mqttp.New(mqttp.ProtocolV311, mqttp.PUBLISH)
notifyMsg, _ := nm.(*mqttp.Publish)
notifyMsg.SetRetain(false)
notifyMsg.SetQoS(mqttp.QoS0) // nolint: errcheck
notifyMsg.SetTopic(t.topic + id + "/connected") // nolint: errcheck
if out, err := json.Marshal(&status); err != nil {
// todo: put reliable message
notifyMsg.SetPayload([]byte("data error"))
} else {
notifyMsg.SetPayload(out)
}
t.topicsManager.Publish(notifyMsg) // nolint: errcheck
t.topicsManager.Retain(notifyMsg) // nolint: errcheck
// notify remove previous disconnect if any
nm, _ = mqttp.New(mqttp.ProtocolV311, mqttp.PUBLISH)
notifyMsg, _ = nm.(*mqttp.Publish)
notifyMsg.SetRetain(false)
notifyMsg.SetQoS(mqttp.QoS0) // nolint: errcheck
notifyMsg.SetTopic(t.topic + id + "/disconnected") // nolint: errcheck
t.topicsManager.Retain(notifyMsg) // nolint: errcheck
}
// Disconnected remove client from statistic
func (t *clients) Disconnected(id string, reason mqttp.ReasonCode) {
atomic.AddUint64(&t.curr.val, ^uint64(0))
nm, _ := mqttp.New(mqttp.ProtocolV311, mqttp.PUBLISH)
notifyMsg, _ := nm.(*mqttp.Publish)
notifyMsg.SetRetain(false)
notifyMsg.SetQoS(mqttp.QoS0) // nolint: errcheck
notifyMsg.SetTopic(t.topic + id + "/disconnected") // nolint: errcheck
notifyPayload := clientDisconnectStatus{
Reason: "normal",
Timestamp: time.Now().Format(time.RFC3339),
}
if out, err := json.Marshal(¬ifyPayload); err != nil {
notifyMsg.SetPayload([]byte("data error"))
} else {
notifyMsg.SetPayload(out)
}
t.topicsManager.Publish(notifyMsg) // nolint: errcheck
// remove connected retained message
nm, _ = mqttp.New(mqttp.ProtocolV311, mqttp.PUBLISH)
notifyMsg, _ = nm.(*mqttp.Publish)
notifyMsg.SetRetain(false)
notifyMsg.SetQoS(mqttp.QoS0) // nolint: errcheck
notifyMsg.SetTopic(t.topic + id + "/connected") // nolint: errcheck
t.topicsManager.Retain(notifyMsg) // nolint: errcheck
}