mirror of
https://github.com/fhmq/hmq.git
synced 2026-08-30 14:24:53 +00:00
add clientID in log for debug
This commit is contained in:
+4
-3
@@ -253,7 +253,7 @@ func (b *Broker) handleConnection(typ int, conn net.Conn, idx uint64) {
|
|||||||
connack.SessionPresent = msg.CleanSession
|
connack.SessionPresent = msg.CleanSession
|
||||||
err = connack.Write(conn)
|
err = connack.Write(conn)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error("send connack error, ", err)
|
log.Error("send connack error, ", err, " clientID = ", msg.ClientIdentifier)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -295,10 +295,11 @@ func (b *Broker) handleConnection(typ int, conn net.Conn, idx uint64) {
|
|||||||
c.mp = msgPool
|
c.mp = msgPool
|
||||||
old, exist = b.clients.Load(cid)
|
old, exist = b.clients.Load(cid)
|
||||||
if exist {
|
if exist {
|
||||||
log.Warn("client exist, close old...")
|
log.Warn("client exist, close old...", " clientID = ", c.info.clientID)
|
||||||
ol, ok := old.(*client)
|
ol, ok := old.(*client)
|
||||||
if ok {
|
if ok {
|
||||||
ol.Close()
|
msg := &Message{client: c, packet: DisconnectdPacket}
|
||||||
|
ol.mp.queue <- msg
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
b.clients.Store(cid, c)
|
b.clients.Store(cid, c)
|
||||||
|
|||||||
+17
-29
@@ -102,7 +102,7 @@ func (c *client) readLoop() {
|
|||||||
}
|
}
|
||||||
packet, err := packets.ReadPacket(nc)
|
packet, err := packets.ReadPacket(nc)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error("read packet error: ", err)
|
log.Error("read packet error: ", err, " clientID = ", c.info.clientID)
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
// log.Info("recv buf: ", packet)
|
// log.Info("recv buf: ", packet)
|
||||||
@@ -124,45 +124,33 @@ func ProcessMessage(msg *Message) {
|
|||||||
if ca == nil {
|
if ca == nil {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
log.Debug("Recv message: ", ca.String(), " clientID = ", c.info.clientID)
|
||||||
switch ca.(type) {
|
switch ca.(type) {
|
||||||
case *packets.ConnackPacket:
|
case *packets.ConnackPacket:
|
||||||
// log.Info("Recv conack message..........")
|
|
||||||
case *packets.ConnectPacket:
|
case *packets.ConnectPacket:
|
||||||
// log.Info("Recv connect message..........")
|
|
||||||
case *packets.PublishPacket:
|
case *packets.PublishPacket:
|
||||||
// log.Info("Recv publish message..........")
|
|
||||||
packet := ca.(*packets.PublishPacket)
|
packet := ca.(*packets.PublishPacket)
|
||||||
c.ProcessPublish(packet)
|
c.ProcessPublish(packet)
|
||||||
case *packets.PubackPacket:
|
case *packets.PubackPacket:
|
||||||
//log.Info("Recv publish ack message..........")
|
|
||||||
case *packets.PubrecPacket:
|
case *packets.PubrecPacket:
|
||||||
//log.Info("Recv publish rec message..........")
|
|
||||||
case *packets.PubrelPacket:
|
case *packets.PubrelPacket:
|
||||||
//log.Info("Recv publish rel message..........")
|
|
||||||
case *packets.PubcompPacket:
|
case *packets.PubcompPacket:
|
||||||
//log.Info("Recv publish ack message..........")
|
|
||||||
case *packets.SubscribePacket:
|
case *packets.SubscribePacket:
|
||||||
// log.Info("Recv subscribe message.....")
|
|
||||||
packet := ca.(*packets.SubscribePacket)
|
packet := ca.(*packets.SubscribePacket)
|
||||||
c.ProcessSubscribe(packet)
|
c.ProcessSubscribe(packet)
|
||||||
case *packets.SubackPacket:
|
case *packets.SubackPacket:
|
||||||
// log.Info("Recv suback message.....")
|
|
||||||
case *packets.UnsubscribePacket:
|
case *packets.UnsubscribePacket:
|
||||||
// log.Info("Recv unsubscribe message.....")
|
|
||||||
packet := ca.(*packets.UnsubscribePacket)
|
packet := ca.(*packets.UnsubscribePacket)
|
||||||
c.ProcessUnSubscribe(packet)
|
c.ProcessUnSubscribe(packet)
|
||||||
case *packets.UnsubackPacket:
|
case *packets.UnsubackPacket:
|
||||||
//log.Info("Recv unsuback message.....")
|
|
||||||
case *packets.PingreqPacket:
|
case *packets.PingreqPacket:
|
||||||
// log.Info("Recv PINGREQ message..........")
|
|
||||||
c.ProcessPing()
|
c.ProcessPing()
|
||||||
case *packets.PingrespPacket:
|
case *packets.PingrespPacket:
|
||||||
//log.Info("Recv PINGRESP message..........")
|
|
||||||
case *packets.DisconnectPacket:
|
case *packets.DisconnectPacket:
|
||||||
// log.Info("Recv DISCONNECT message.......")
|
|
||||||
c.Close()
|
c.Close()
|
||||||
default:
|
default:
|
||||||
log.Info("Recv Unknow message.......")
|
log.Info("Recv Unknow message.......", " clientID = ", c.info.clientID)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -173,7 +161,7 @@ func (c *client) ProcessPublish(packet *packets.PublishPacket) {
|
|||||||
|
|
||||||
topic := packet.TopicName
|
topic := packet.TopicName
|
||||||
if !c.CheckTopicAuth(PUB, topic) {
|
if !c.CheckTopicAuth(PUB, topic) {
|
||||||
log.Error("Pub Topics Auth failed, ", topic)
|
log.Error("Pub Topics Auth failed, ", topic, " clientID = ", c.info.clientID)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -184,21 +172,21 @@ func (c *client) ProcessPublish(packet *packets.PublishPacket) {
|
|||||||
puback := packets.NewControlPacket(packets.Puback).(*packets.PubackPacket)
|
puback := packets.NewControlPacket(packets.Puback).(*packets.PubackPacket)
|
||||||
puback.MessageID = packet.MessageID
|
puback.MessageID = packet.MessageID
|
||||||
if err := c.WriterPacket(puback); err != nil {
|
if err := c.WriterPacket(puback); err != nil {
|
||||||
log.Error("send puback error, ", err)
|
log.Error("send puback error, ", err, " clientID = ", c.info.clientID)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
c.ProcessPublishMessage(packet)
|
c.ProcessPublishMessage(packet)
|
||||||
case QosExactlyOnce:
|
case QosExactlyOnce:
|
||||||
return
|
return
|
||||||
default:
|
default:
|
||||||
log.Error("publish with unknown qos")
|
log.Error("publish with unknown qos", " clientID = ", c.info.clientID)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
if packet.Retain {
|
if packet.Retain {
|
||||||
if b := c.broker; b != nil {
|
if b := c.broker; b != nil {
|
||||||
err := b.rl.Insert(topic, packet)
|
err := b.rl.Insert(topic, packet)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error("Insert Retain Message error: ", err)
|
log.Error("Insert Retain Message error: ", err, " clientID = ", c.info.clientID)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -232,7 +220,7 @@ func (c *client) ProcessPublishMessage(packet *packets.PublishPacket) {
|
|||||||
if sub != nil {
|
if sub != nil {
|
||||||
err := sub.client.WriterPacket(packet)
|
err := sub.client.WriterPacket(packet)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error("process message for psub error, ", err)
|
log.Error("process message for psub error, ", err, " clientID = ", c.info.clientID)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -258,7 +246,7 @@ func (c *client) ProcessPublishMessage(packet *packets.PublishPacket) {
|
|||||||
if sub != nil {
|
if sub != nil {
|
||||||
err := sub.client.WriterPacket(packet)
|
err := sub.client.WriterPacket(packet)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error("send publish error, ", err)
|
log.Error("send publish error, ", err, " clientID = ", c.info.clientID)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -312,7 +300,7 @@ func (c *client) ProcessSubscribe(packet *packets.SubscribePacket) {
|
|||||||
t := topic
|
t := topic
|
||||||
//check topic auth for client
|
//check topic auth for client
|
||||||
if !c.CheckTopicAuth(SUB, topic) {
|
if !c.CheckTopicAuth(SUB, topic) {
|
||||||
log.Error("Sub topic Auth failed: ", topic)
|
log.Error("Sub topic Auth failed: ", topic, " clientID = ", c.info.clientID)
|
||||||
retcodes = append(retcodes, QosFailure)
|
retcodes = append(retcodes, QosFailure)
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
@@ -359,7 +347,7 @@ func (c *client) ProcessSubscribe(packet *packets.SubscribePacket) {
|
|||||||
}
|
}
|
||||||
err := b.sl.Insert(sub)
|
err := b.sl.Insert(sub)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error("Insert subscription error: ", err)
|
log.Error("Insert subscription error: ", err, " clientID = ", c.info.clientID)
|
||||||
retcodes = append(retcodes, QosFailure)
|
retcodes = append(retcodes, QosFailure)
|
||||||
} else {
|
} else {
|
||||||
retcodes = append(retcodes, qoss[i])
|
retcodes = append(retcodes, qoss[i])
|
||||||
@@ -369,7 +357,7 @@ func (c *client) ProcessSubscribe(packet *packets.SubscribePacket) {
|
|||||||
|
|
||||||
err := c.WriterPacket(suback)
|
err := c.WriterPacket(suback)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error("send suback error, ", err)
|
log.Error("send suback error, ", err, " clientID = ", c.info.clientID)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
//broadcast subscribe message
|
//broadcast subscribe message
|
||||||
@@ -381,7 +369,7 @@ func (c *client) ProcessSubscribe(packet *packets.SubscribePacket) {
|
|||||||
for _, t := range topics {
|
for _, t := range topics {
|
||||||
packets := b.rl.Match(t)
|
packets := b.rl.Match(t)
|
||||||
for _, packet := range packets {
|
for _, packet := range packets {
|
||||||
log.Info("process retain message: ", packet)
|
log.Info("process retain message: ", packet, " clientID = ", c.info.clientID)
|
||||||
if packet != nil {
|
if packet != nil {
|
||||||
c.WriterPacket(packet)
|
c.WriterPacket(packet)
|
||||||
}
|
}
|
||||||
@@ -432,7 +420,7 @@ func (c *client) ProcessUnSubscribe(packet *packets.UnsubscribePacket) {
|
|||||||
|
|
||||||
err := c.WriterPacket(unsuback)
|
err := c.WriterPacket(unsuback)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error("send unsuback error, ", err)
|
log.Error("send unsuback error, ", err, " clientID = ", c.info.clientID)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
// //process ubsubscribe message
|
// //process ubsubscribe message
|
||||||
@@ -461,7 +449,7 @@ func (c *client) ProcessPing() {
|
|||||||
resp := packets.NewControlPacket(packets.Pingresp).(*packets.PingrespPacket)
|
resp := packets.NewControlPacket(packets.Pingresp).(*packets.PingrespPacket)
|
||||||
err := c.WriterPacket(resp)
|
err := c.WriterPacket(resp)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error("send PingResponse error, ", err)
|
log.Error("send PingResponse error, ", err, " clientID = ", c.info.clientID)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -489,7 +477,7 @@ func (c *client) Close() {
|
|||||||
for _, sub := range subs {
|
for _, sub := range subs {
|
||||||
err := b.sl.Remove(sub)
|
err := b.sl.Remove(sub)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error("closed client but remove sublist error, ", err)
|
log.Error("closed client but remove sublist error, ", err, " clientID = ", c.info.clientID)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if c.typ == CLIENT {
|
if c.typ == CLIENT {
|
||||||
|
|||||||
Reference in New Issue
Block a user