代码拉取完成,页面将自动刷新
同步操作将从 JUMEI_ARCH/volantmq 强制同步,此操作会覆盖自 Fork 仓库以来所做的任何修改,且无法恢复!!!
确定后同步将在后台操作,完成时将刷新页面,请耐心等待。
package systree
import (
"encoding/json"
"strconv"
"sync/atomic"
"github.com/VolantMQ/volantmq/packet"
"github.com/VolantMQ/volantmq/types"
"github.com/prometheus/client_golang/prometheus"
)
// SessionCreatedStatus report when session status once created
type SessionCreatedStatus struct {
ExpiryInterval string `json:"expiryInterval,omitempty"`
WillDelay string `json:"willDelay,omitempty"`
Timestamp string `json:"timestamp"`
Clean bool `json:"clean"`
}
// SessionDeletedStatus report when session status once deleted
type SessionDeletedStatus struct {
Timestamp string `json:"timestamp"`
Reason string `json:"reason"`
}
type sessions struct {
stat
topicsManager types.TopicMessenger
topic string
prom prometheus.Collector
}
type promSessions struct {
cur *promMetricEntry
max *promMetricEntry
}
func (p *promSessions) Describe(ch chan<- *prometheus.Desc) {
ch <- p.cur.desc
ch <- p.max.desc
}
func (p *promSessions) Collect(ch chan<- prometheus.Metric) {
val := p.cur.valFunc()
v, err := strconv.ParseFloat(string(val), 64)
if err != nil {
ch <- prometheus.MustNewConstMetric(p.cur.desc, p.cur.valType, 0)
} else {
ch <- prometheus.MustNewConstMetric(p.cur.desc, p.cur.valType, v)
}
val = p.max.valFunc()
v, err = strconv.ParseFloat(string(val), 64)
if err != nil {
ch <- prometheus.MustNewConstMetric(p.max.desc, p.max.valType, 0)
} else {
ch <- prometheus.MustNewConstMetric(p.max.desc, p.max.valType, v)
}
}
func newSessionsCollector(t stat) prometheus.Collector {
p := &promSessions{
cur: &promMetricEntry{
valFunc: t.curr.getValue,
desc: prometheus.NewDesc("jmmqtt_cur_sessions",
"Number of current sessions",
nil, nil,
),
valType: prometheus.CounterValue,
},
max: &promMetricEntry{
valFunc: t.max.getValue,
desc: prometheus.NewDesc("jmmqtt_max_sessions",
"Number of max sessions",
nil, nil,
),
valType: prometheus.CounterValue,
},
}
return p
}
func newSessions(topicPrefix string, retained *[]types.RetainObject) sessions {
c := sessions{
stat: newStat(topicPrefix+"/stats/sessions", retained),
topic: topicPrefix + "/sessions/",
}
c.prom = newSessionsCollector(c.stat)
prometheus.MustRegister(c.prom)
return c
}
// Connected add to statistic new client
func (t *sessions) Created(id string, status *SessionCreatedStatus) {
newVal := atomic.AddUint64(&t.curr.val, 1)
if atomic.LoadUint64(&t.max.val) < newVal {
atomic.StoreUint64(&t.max.val, newVal)
}
if t.topicsManager != nil {
// notify client connected
nm, _ := packet.New(packet.ProtocolV311, packet.PUBLISH)
notifyMsg, _ := nm.(*packet.Publish)
notifyMsg.SetRetain(false)
notifyMsg.SetQoS(packet.QoS0) // nolint: errcheck
notifyMsg.SetTopic(t.topic + id + "/created") // 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
}
}
// Disconnected remove client from statistic
func (t *sessions) Removed(id string, status *SessionDeletedStatus) {
atomic.AddUint64(&t.curr.val, ^uint64(0))
if t.topicsManager != nil {
nm, _ := packet.New(packet.ProtocolV311, packet.PUBLISH)
notifyMsg, _ := nm.(*packet.Publish)
notifyMsg.SetRetain(false)
notifyMsg.SetQoS(packet.QoS0) // nolint: errcheck
notifyMsg.SetTopic(t.topic + id) // nolint: errcheck
t.topicsManager.Retain(notifyMsg) // nolint: errcheck
nm, _ = packet.New(packet.ProtocolV311, packet.PUBLISH)
notifyMsg, _ = nm.(*packet.Publish)
notifyMsg.SetRetain(false)
notifyMsg.SetQoS(packet.QoS0) // nolint: errcheck
notifyMsg.SetTopic(t.topic + id) // nolint: errcheck
notifyMsg.SetTopic(t.topic + id + "/removed") // nolint: errcheck
if out, err := json.Marshal(&status); err != nil {
notifyMsg.SetPayload([]byte("data error"))
} else {
notifyMsg.SetPayload(out)
}
t.topicsManager.Publish(notifyMsg) // nolint: errcheck
}
}
此处可能存在不合适展示的内容,页面不予展示。您可通过相关编辑功能自查并修改。
如您确认内容无涉及 不当用语 / 纯广告导流 / 暴力 / 低俗色情 / 侵权 / 盗版 / 虚假 / 无价值内容或违法国家有关法律法规的内容,可点击提交进行申诉,我们将尽快为您处理。