Class MqttDecoder

All Implemented Interfaces:
ChannelHandler, ChannelInboundHandler

public final class MqttDecoder extends ReplayingDecoder<MqttDecoder.DecoderState>
Decodes Mqtt messages from bytes, following the MQTT protocol specification v3.1 or v5.0, depending on the version specified in the CONNECT message that first goes through the channel.
  • Field Details

    • mqttFixedHeader

      private MqttFixedHeader mqttFixedHeader
    • variableHeader

      private Object variableHeader
    • bytesRemainingInVariablePart

      private int bytesRemainingInVariablePart
    • maxBytesInMessage

      private final int maxBytesInMessage
    • maxClientIdLength

      private final int maxClientIdLength
    • strictUtf8Validation

      private final boolean strictUtf8Validation
    • utf8Decoder

      private CharsetDecoder utf8Decoder
    • maxAllowedRemainingBytes

      private int maxAllowedRemainingBytes
  • Constructor Details

    • MqttDecoder

      public MqttDecoder()
    • MqttDecoder

      public MqttDecoder(int maxBytesInMessage)
    • MqttDecoder

      public MqttDecoder(int maxBytesInMessage, int maxClientIdLength)
    • MqttDecoder

      public MqttDecoder(int maxBytesInMessage, int maxClientIdLength, boolean strictUtf8Validation)
      Creates a new MqttDecoder.
      Parameters:
      maxBytesInMessage - the maximum number of bytes a decoded message may consume.
      maxClientIdLength - the maximum length of the Client Identifier (CONNECT payload).
      strictUtf8Validation - if true (default), every UTF-8 Encoded String is validated according to MQTT 3.1.1 and MQTT 5.0 malformed UTF-8 sequences (including surrogates and overlong forms) and an embedded U+0000 are rejected as a Malformed Packet. If false, the legacy behaviour is preserved, malformed bytes are silently replaced with U+FFFD and U+0000 is accepted.
  • Method Details

    • decode

      protected void decode(ChannelHandlerContext ctx, ByteBuf buffer, List<Object> out) throws Exception
      Description copied from class: ByteToMessageDecoder
      Decode the from one ByteBuf to an other. This method will be called till either the input ByteBuf has nothing to read when return from this method or till nothing was read from the input ByteBuf.
      Specified by:
      decode in class ByteToMessageDecoder
      Parameters:
      ctx - the ChannelHandlerContext which this ByteToMessageDecoder belongs to
      buffer - the ByteBuf from which to read data
      out - the List to which decoded messages should be added
      Throws:
      Exception - is thrown if an error occurs
    • invalidMessage

      private MqttMessage invalidMessage(Throwable cause)
    • checkMaxMessageLengthRemaining

      private void checkMaxMessageLengthRemaining(int maxAllowedRemainingBytes)
    • decodeFixedHeader

      private MqttFixedHeader decodeFixedHeader(ChannelHandlerContext ctx, ByteBuf buffer, int maxAllowedRemainingBytes)
      Decodes the fixed header. It's one byte for the flags and then variable bytes for the remaining length.
      Parameters:
      buffer - the buffer to decode from
      Returns:
      the fixed header
    • parseRemainingLength

      private int parseRemainingLength(ByteBuf buffer, MqttMessageType messageType, int maxAllowedRemainingBytes)
    • decodeVariableHeader

      private Object decodeVariableHeader(ChannelHandlerContext ctx, ByteBuf buffer, MqttFixedHeader mqttFixedHeader, int maxAllowedRemainingBytes)
      Decodes the variable header (if any)
      Parameters:
      buffer - the buffer to decode from
      mqttFixedHeader - MqttFixedHeader of the same message
      maxAllowedRemainingBytes - the maximum number of bytes permitted to remain
      Returns:
      the variable header
    • decodeConnectionVariableHeader

      private MqttConnectVariableHeader decodeConnectionVariableHeader(ChannelHandlerContext ctx, ByteBuf buffer, int maxAllowedRemainingBytes)
    • decodeConnAckVariableHeader

      private MqttConnAckVariableHeader decodeConnAckVariableHeader(ChannelHandlerContext ctx, ByteBuf buffer, int maxAllowedRemainingBytes)
    • decodeMessageIdAndPropertiesVariableHeader

      private MqttMessageIdAndPropertiesVariableHeader decodeMessageIdAndPropertiesVariableHeader(ChannelHandlerContext ctx, ByteBuf buffer, int maxAllowedRemainingBytes)
    • decodePubReplyMessage

      private MqttPubReplyMessageVariableHeader decodePubReplyMessage(ByteBuf buffer, int maxAllowedRemainingBytes)
    • decodeReasonCodeAndPropertiesVariableHeader

      private MqttReasonCodeAndPropertiesVariableHeader decodeReasonCodeAndPropertiesVariableHeader(ByteBuf buffer, int maxAllowedRemainingBytes)
    • decodePublishVariableHeader

      private MqttPublishVariableHeader decodePublishVariableHeader(ChannelHandlerContext ctx, ByteBuf buffer, MqttFixedHeader mqttFixedHeader, int maxAllowedRemainingBytes)
    • decodeMessageId

      private static int decodeMessageId(ByteBuf buffer)
      Returns:
      messageId with numberOfBytesConsumed is 2
    • decodePayload

      private Object decodePayload(ByteBuf buffer, MqttMessageType messageType, int maxClientIdLength, Object variableHeader, int maxAllowedRemainingBytes)
      Decodes the payload.
      Parameters:
      buffer - the buffer to decode from
      messageType - type of the message being decoded
      variableHeader - variable header of the same message
      Returns:
      the payload
    • decodeConnectionPayload

      private MqttConnectPayload decodeConnectionPayload(ByteBuf buffer, int maxClientIdLength, MqttConnectVariableHeader mqttConnectVariableHeader, int maxAllowedRemainingBytes)
    • decodeSubscribePayload

      private MqttSubscribePayload decodeSubscribePayload(ByteBuf buffer, int maxAllowedRemainingBytes)
    • decodeSubackPayload

      private MqttSubAckPayload decodeSubackPayload(ByteBuf buffer, int maxAllowedRemainingBytes)
    • decodeUnsubAckPayload

      private MqttUnsubAckPayload decodeUnsubAckPayload(ByteBuf buffer, int maxAllowedRemainingBytes)
    • decodeUnsubscribePayload

      private MqttUnsubscribePayload decodeUnsubscribePayload(ByteBuf buffer, int maxAllowedRemainingBytes)
    • decodePublishPayload

      private ByteBuf decodePublishPayload(ByteBuf buffer, int maxAllowedRemainingBytes)
    • validateNoBytesRemain

      private void validateNoBytesRemain(int numberOfBytesConsumed)
    • decodeString

      private MqttDecoder.Result<String> decodeString(ByteBuf buffer, int maxAllowedRemainingBytes)
    • decodeString

      private MqttDecoder.Result<String> decodeString(ByteBuf buffer, int minBytes, int maxBytes, int maxAllowedRemainingBytes)
    • readStrictUtf8

      private String readStrictUtf8(ByteBuf buffer, int length)
      Reads length bytes from buffer and decodes them as a strictly validated UTF-8 Encoded String per MQTT 3.1.1 and MQTT 5.0. Throws a DecoderException if the sequence is malformed or contains U+0000.
    • decodeByteArray

      private byte[] decodeByteArray(ByteBuf buffer, int maxAllowedRemainingBytes)
      Returns:
      the decoded byte[], numberOfBytesConsumed = byte[].length + 2
    • packInts

      private static long packInts(int a, int b)
    • unpackA

      private static int unpackA(long ints)
    • unpackB

      private static int unpackB(long ints)
    • decodeMsbLsb

      private static int decodeMsbLsb(ByteBuf buffer)
      numberOfBytesConsumed = 2. return decoded result.
    • decodeVariableByteInteger

      private long decodeVariableByteInteger(ByteBuf buffer, int maxAllowedRemainingBytes)
      See 1.5.5 Variable Byte Integer section of MQTT 5.0 specification for encoding/decoding rules
      Parameters:
      buffer - the buffer to decode from
      Returns:
      result pack with a = decoded integer, b = numberOfBytesConsumed. Need to unpack to read them.
      Throws:
      DecoderException - if bad MQTT protocol limits Remaining Length
    • decodeProperties

      private MqttDecoder.Result<MqttProperties> decodeProperties(ByteBuf buffer, int maxAllowedRemainingBytes)