Problems with MQTT server sending data
- Dominant language
- C#
- Stars
- 4.3k
- Forks
- 1k
- PR merge metrics
- No merged PRs in 30d
Description
After my device sends some basic environment information of the device to the server, the mqtt server pings the data sent by the client before.
this sample code
``` c#
class MqttHandler : SimpleChannelInboundHandler
{
public MqttHandler() { }
///
///
///
///
///
private async Task ProcessMqttMsgAsync(IChannelHandlerContext context, Packet msg)
{
switch (msg.PacketType)
{
case PacketType.CONNECT:
await context.WriteAndFlushAsync(new ConnAckPacket { ReturnCode = ConnectReturnCode.Accepted, SessionPresent = true });
break;
case PacketType.PUBACK:
if (msg is PubAckPacket pubAckPacket)
{
var msgId = pubAckPacket.PacketId;
}
break;
case PacketType.PUBCOMP:
if (msg is PubCompPacket pubCompPacket)
{
var msgId = pubCompPacket.PacketId;
}
break;
case PacketType.PUBREC:
if (msg is PubRecPacket pubRecPacket)
{
}
break;
case PacketType.PUBREL:
if (msg is PubRelPacket pubRelPacket)
{
}
break;
case PacketType.SUBSCRIBE:
await context.WriteAndFlushAsync(SubAckPacket.InResponseTo(msg as SubscribePacket, QualityOfService.ExactlyOnce));
break;
case PacketType.UNSUBSCRIBE:
await context.WriteAndFlushAsync(UnsubAckPacket.InResponseTo(msg as UnsubscribePacket));
break;
case PacketType.PINGREQ:
await context.WriteAndFlushAsync(PingRespPacket.Instance);
break;
case PacketType.DISCONNECT:
break;
default:
break;
}
}
private async Task ProcessConnectAsync(IChannelHandlerContext context, ConnectPacket connectPacket)
{
await context.WriteAndFlushAsync(CreateMqttConnectionAck(ConnectReturnCode.Accepted, connectPacket));
}
private ConnAckPacket CreateMqttConnectionAck(ConnectReturnCode returnCode, ConnectPacket connectPacket)
{
if (connectPacket == null)
return null;
return new ConnAckPacket { ReturnCode = returnCode, SessionPresent = !connectPacket.CleanSession };
}
protected override void ChannelRead0(IChannelHandlerContext context, Packet message)
{
if (message is PublishPacket publishPacket)
{
Console.WriteLine(publishPacket.Payload.ToString(encoding: Encoding.UTF8));
return;
}
try
{
if (message is Packet msg)
{
ProcessMqttMsgAsync(context, msg).GetAwaiter();
}
else
{
context.CloseAsync();
}
}
finally
{
}
}
}
```
``` c#
public static void Main() => RunServerAsync().Wait();
public static async Task RunServerAsync()
{
IChannel boundChannel;
var bossGroup = new MultithreadEventLoopGroup(1);
var workerGroup = new MultithreadEventLoopGroup();
try
{
var bootstrap = new ServerBootstrap();
bootstrap.Group(bossGroup, workerGroup);
bootstrap.Channel();
bootstrap.Option(ChannelOption.SoBacklog, 100)
.Option(ChannelOption.SoKeepalive, false)
.ChildHandler(new ActionChannelInitializer(channel =>
{
IChannelPipeline pipeline = channel.Pipeline;
pipeline.AddLast(MqttEncoder.Instance, new MqttDecoder(true, 64 * 1024), new MqttHandler());
}
));
boundChannel = await bootstrap.BindAsync(1883);
Console.WriteLine("mqtt server is start");
Console.ReadLine();
}
catch(Exception ex)
{
Console.WriteLine("mqtt server failue");
}
}
}
```
Contributor guide
Assessment
This issue has not been assessed yet.