View Javadoc
1   /*
2    * Copyright 2014 The Netty Project
3    *
4    * The Netty Project licenses this file to you under the Apache License,
5    * version 2.0 (the "License"); you may not use this file except in compliance
6    * with the License. You may obtain a copy of the License at:
7    *
8    *   https://www.apache.org/licenses/LICENSE-2.0
9    *
10   * Unless required by applicable law or agreed to in writing, software
11   * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
12   * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
13   * License for the specific language governing permissions and limitations
14   * under the License.
15   */
16  
17  package io.netty.handler.codec.mqtt;
18  
19  import io.netty.buffer.ByteBuf;
20  import io.netty.buffer.ByteBufAllocator;
21  import io.netty.buffer.Unpooled;
22  import io.netty.channel.ChannelHandler;
23  import io.netty.channel.ChannelHandlerContext;
24  import io.netty.handler.codec.EncoderException;
25  import io.netty.handler.codec.MessageToMessageEncoder;
26  import io.netty.util.internal.EmptyArrays;
27  
28  import java.util.List;
29  
30  import static io.netty.buffer.ByteBufUtil.*;
31  import static io.netty.handler.codec.mqtt.MqttCodecUtil.getMqttVersion;
32  import static io.netty.handler.codec.mqtt.MqttCodecUtil.isValidClientId;
33  import static io.netty.handler.codec.mqtt.MqttCodecUtil.isValidPublishTopicName;
34  import static io.netty.handler.codec.mqtt.MqttCodecUtil.isValidUserName;
35  import static io.netty.handler.codec.mqtt.MqttCodecUtil.setMqttVersion;
36  import static io.netty.handler.codec.mqtt.MqttConstant.DEFAULT_MAX_CLIENT_ID_LENGTH;
37  
38  /**
39   * Encodes Mqtt messages into bytes following the protocol specification v3.1
40   * as described here <a href="https://public.dhe.ibm.com/software/dw/webservices/ws-mqtt/mqtt-v3r1.html">MQTTV3.1</a>
41   * or v5.0 as described here <a href="https://docs.oasis-open.org/mqtt/mqtt/v5.0/mqtt-v5.0.html">MQTTv5.0</a> -
42   * depending on the version specified in the first CONNECT message that goes through the channel.
43   */
44  @ChannelHandler.Sharable
45  public final class MqttEncoder extends MessageToMessageEncoder<MqttMessage> {
46  
47      public static final MqttEncoder INSTANCE = new MqttEncoder();
48  
49      private MqttEncoder() { }
50  
51      @Override
52      protected void encode(ChannelHandlerContext ctx, MqttMessage msg, List<Object> out) throws Exception {
53          out.add(doEncode(ctx, msg));
54      }
55  
56      /**
57       * This is the main encoding method.
58       * It's only visible for testing.
59       *
60       * @param message MQTT message to encode
61       * @return ByteBuf with encoded bytes
62       */
63      static ByteBuf doEncode(ChannelHandlerContext ctx,
64                       MqttMessage message) {
65  
66          switch (message.fixedHeader().messageType()) {
67              case CONNECT:
68                  return encodeConnectMessage(ctx, (MqttConnectMessage) message);
69  
70              case CONNACK:
71                  return encodeConnAckMessage(ctx, (MqttConnAckMessage) message);
72  
73              case PUBLISH:
74                  return encodePublishMessage(ctx, (MqttPublishMessage) message);
75  
76              case SUBSCRIBE:
77                  return encodeSubscribeMessage(ctx, (MqttSubscribeMessage) message);
78  
79              case UNSUBSCRIBE:
80                  return encodeUnsubscribeMessage(ctx,  (MqttUnsubscribeMessage) message);
81  
82              case SUBACK:
83                  return encodeSubAckMessage(ctx, (MqttSubAckMessage) message);
84  
85              case UNSUBACK:
86                  if (message instanceof MqttUnsubAckMessage) {
87                      return encodeUnsubAckMessage(ctx, (MqttUnsubAckMessage) message);
88                  }
89                  return encodeMessageWithOnlySingleByteFixedHeaderAndMessageId(ctx.alloc(), message);
90  
91              case PUBACK:
92              case PUBREC:
93              case PUBREL:
94              case PUBCOMP:
95                  return encodePubReplyMessage(ctx, message);
96  
97              case DISCONNECT:
98              case AUTH:
99                  return encodeReasonCodePlusPropertiesMessage(ctx, message);
100 
101             case PINGREQ:
102             case PINGRESP:
103                 return encodeMessageWithOnlySingleByteFixedHeader(ctx.alloc(), message);
104 
105             default:
106                 throw new IllegalArgumentException(
107                         "Unknown message type: " + message.fixedHeader().messageType().value());
108         }
109     }
110 
111     private static ByteBuf encodeConnectMessage(
112             ChannelHandlerContext ctx,
113             MqttConnectMessage message) {
114         int payloadBufferSize = 0;
115 
116         MqttFixedHeader mqttFixedHeader = message.fixedHeader();
117         MqttConnectVariableHeader variableHeader = message.variableHeader();
118         MqttConnectPayload payload = message.payload();
119         MqttVersion mqttVersion = MqttVersion.fromProtocolNameAndLevel(variableHeader.name(),
120                 (byte) variableHeader.version());
121         setMqttVersion(ctx, mqttVersion);
122 
123         // MQTT 3.1 and 3.1.1 require the Password Flag to be 0 when the User Name Flag is 0.
124         if ((mqttVersion == MqttVersion.MQTT_3_1 || mqttVersion == MqttVersion.MQTT_3_1_1) &&
125                 !variableHeader.hasUserName() && variableHeader.hasPassword()) {
126             throw new EncoderException("Without a username, the password MUST be not set");
127         }
128 
129         // Client id
130         String clientIdentifier = payload.clientIdentifier();
131         if (!isValidClientId(mqttVersion, DEFAULT_MAX_CLIENT_ID_LENGTH, clientIdentifier, false)) {
132             throw new MqttIdentifierRejectedException("invalid clientIdentifier");
133         }
134         int clientIdentifierBytes = utf8Bytes(clientIdentifier);
135         payloadBufferSize += 2 + clientIdentifierBytes;
136 
137         // Will topic and message
138         String willTopic = payload.willTopic();
139         int willTopicBytes = nullableUtf8Bytes(willTopic);
140         byte[] willMessage = payload.willMessageInBytes();
141         byte[] willMessageBytes = willMessage != null ? willMessage : EmptyArrays.EMPTY_BYTES;
142         if (variableHeader.isWillFlag()) {
143             if (!isValidPublishTopicName(willTopic)) {
144                 throw new MqttIdentifierRejectedException("invalid willTopic");
145             }
146             payloadBufferSize += 2 + willTopicBytes;
147             payloadBufferSize += 2 + willMessageBytes.length;
148         }
149 
150         String userName = payload.userName();
151         int userNameBytes = nullableUtf8Bytes(userName);
152         if (variableHeader.hasUserName()) {
153             if (!isValidUserName(userName)) {
154                 throw new MqttIdentifierRejectedException("invalid userName");
155             }
156             payloadBufferSize += 2 + userNameBytes;
157         }
158 
159         byte[] password = payload.passwordInBytes();
160         byte[] passwordBytes = password != null ? password : EmptyArrays.EMPTY_BYTES;
161         if (variableHeader.hasPassword()) {
162             payloadBufferSize += 2 + passwordBytes.length;
163         }
164 
165         // Fixed and variable header
166         byte[] protocolNameBytes = mqttVersion.protocolNameBytes();
167         ByteBuf propertiesBuf = encodePropertiesIfNeeded(
168                 mqttVersion,
169                 ctx.alloc(),
170                 message.variableHeader().properties());
171         try {
172             final ByteBuf willPropertiesBuf;
173             if (variableHeader.isWillFlag()) {
174                 willPropertiesBuf = encodePropertiesIfNeeded(mqttVersion, ctx.alloc(), payload.willProperties());
175                 payloadBufferSize += willPropertiesBuf.readableBytes();
176             } else {
177                 willPropertiesBuf = Unpooled.EMPTY_BUFFER;
178             }
179             try {
180                 int variableHeaderBufferSize = 2 + protocolNameBytes.length + 4 + propertiesBuf.readableBytes();
181 
182                 int variablePartSize = variableHeaderBufferSize + payloadBufferSize;
183                 int fixedHeaderBufferSize = 1 + getVariableLengthInt(variablePartSize);
184                 ByteBuf buf = ctx.alloc().buffer(fixedHeaderBufferSize + variablePartSize);
185                 buf.writeByte(getFixedHeaderByte1(mqttFixedHeader));
186                 writeVariableLengthInt(buf, variablePartSize);
187 
188                 buf.writeShort(protocolNameBytes.length);
189                 buf.writeBytes(protocolNameBytes);
190 
191                 buf.writeByte(variableHeader.version());
192                 buf.writeByte(getConnVariableHeaderFlag(variableHeader));
193                 buf.writeShort(variableHeader.keepAliveTimeSeconds());
194                 buf.writeBytes(propertiesBuf);
195 
196                 // Payload
197                 writeExactUTF8String(buf, clientIdentifier, clientIdentifierBytes);
198                 if (variableHeader.isWillFlag()) {
199                     buf.writeBytes(willPropertiesBuf);
200                     writeExactUTF8String(buf, willTopic, willTopicBytes);
201                     buf.writeShort(willMessageBytes.length);
202                     buf.writeBytes(willMessageBytes, 0, willMessageBytes.length);
203                 }
204                 if (variableHeader.hasUserName()) {
205                     writeExactUTF8String(buf, userName, userNameBytes);
206                 }
207                 if (variableHeader.hasPassword()) {
208                     buf.writeShort(passwordBytes.length);
209                     buf.writeBytes(passwordBytes, 0, passwordBytes.length);
210                 }
211                 return buf;
212             } finally {
213                 willPropertiesBuf.release();
214             }
215         } finally {
216             propertiesBuf.release();
217         }
218     }
219 
220     private static int getConnVariableHeaderFlag(MqttConnectVariableHeader variableHeader) {
221         int flagByte = 0;
222         if (variableHeader.hasUserName()) {
223             flagByte |= 0x80;
224         }
225         if (variableHeader.hasPassword()) {
226             flagByte |= 0x40;
227         }
228         if (variableHeader.isWillRetain()) {
229             flagByte |= 0x20;
230         }
231         flagByte |= (variableHeader.willQos() & 0x03) << 3;
232         if (variableHeader.isWillFlag()) {
233             flagByte |= 0x04;
234         }
235         if (variableHeader.isCleanSession()) {
236             flagByte |= 0x02;
237         }
238         return flagByte;
239     }
240 
241     private static ByteBuf encodeConnAckMessage(
242             ChannelHandlerContext ctx,
243             MqttConnAckMessage message) {
244         final MqttVersion mqttVersion = getMqttVersion(ctx);
245         ByteBuf propertiesBuf = encodePropertiesIfNeeded(mqttVersion,
246                 ctx.alloc(),
247                 message.variableHeader().properties());
248 
249         try {
250             ByteBuf buf = ctx.alloc().buffer(4 + propertiesBuf.readableBytes());
251             buf.writeByte(getFixedHeaderByte1(message.fixedHeader()));
252             writeVariableLengthInt(buf, 2 + propertiesBuf.readableBytes());
253             buf.writeByte(message.variableHeader().isSessionPresent() ? 0x01 : 0x00);
254             buf.writeByte(message.variableHeader().connectReturnCode().byteValue());
255             buf.writeBytes(propertiesBuf);
256             return buf;
257         } finally {
258             propertiesBuf.release();
259         }
260     }
261 
262     private static ByteBuf encodeSubscribeMessage(
263             ChannelHandlerContext ctx,
264             MqttSubscribeMessage message) {
265         MqttVersion mqttVersion = getMqttVersion(ctx);
266         ByteBuf propertiesBuf = encodePropertiesIfNeeded(mqttVersion,
267                 ctx.alloc(),
268                 message.idAndPropertiesVariableHeader().properties());
269 
270         try {
271             final int variableHeaderBufferSize = 2 + propertiesBuf.readableBytes();
272             int payloadBufferSize = 0;
273 
274             MqttFixedHeader mqttFixedHeader = message.fixedHeader();
275             MqttMessageIdVariableHeader variableHeader = message.variableHeader();
276             MqttSubscribePayload payload = message.payload();
277 
278             for (MqttTopicSubscription topic : payload.topicSubscriptions()) {
279                 String topicName = topic.topicName();
280                 int topicNameBytes = utf8Bytes(topicName);
281                 payloadBufferSize += 2 + topicNameBytes;
282                 payloadBufferSize += 1;
283             }
284 
285             int variablePartSize = variableHeaderBufferSize + payloadBufferSize;
286             int fixedHeaderBufferSize = 1 + getVariableLengthInt(variablePartSize);
287 
288             ByteBuf buf = ctx.alloc().buffer(fixedHeaderBufferSize + variablePartSize);
289             buf.writeByte(getFixedHeaderByte1(mqttFixedHeader));
290             writeVariableLengthInt(buf, variablePartSize);
291 
292             // Variable Header
293             int messageId = variableHeader.messageId();
294             buf.writeShort(messageId);
295             buf.writeBytes(propertiesBuf);
296 
297             // Payload
298             for (MqttTopicSubscription topic : payload.topicSubscriptions()) {
299                 writeEagerUTF8String(buf, topic.topicName());
300                 if (mqttVersion == MqttVersion.MQTT_3_1_1 || mqttVersion == MqttVersion.MQTT_3_1) {
301                     buf.writeByte(topic.qualityOfService().value());
302                 } else {
303                     final MqttSubscriptionOption option = topic.option();
304 
305                     int optionEncoded = option.retainHandling().value() << 4;
306                     if (option.isRetainAsPublished()) {
307                         optionEncoded |= 0x08;
308                     }
309                     if (option.isNoLocal()) {
310                         optionEncoded |= 0x04;
311                     }
312                     optionEncoded |= option.qos().value();
313 
314                     buf.writeByte(optionEncoded);
315                 }
316             }
317 
318             return buf;
319         } finally {
320             propertiesBuf.release();
321         }
322     }
323 
324     private static ByteBuf encodeUnsubscribeMessage(
325             ChannelHandlerContext ctx,
326             MqttUnsubscribeMessage message) {
327         MqttVersion mqttVersion = getMqttVersion(ctx);
328         ByteBuf propertiesBuf = encodePropertiesIfNeeded(mqttVersion,
329                 ctx.alloc(),
330                 message.idAndPropertiesVariableHeader().properties());
331 
332         try {
333             final int variableHeaderBufferSize = 2 + propertiesBuf.readableBytes();
334             int payloadBufferSize = 0;
335 
336             MqttFixedHeader mqttFixedHeader = message.fixedHeader();
337             MqttMessageIdVariableHeader variableHeader = message.variableHeader();
338             MqttUnsubscribePayload payload = message.payload();
339 
340             for (String topicName : payload.topics()) {
341                 int topicNameBytes = utf8Bytes(topicName);
342                 payloadBufferSize += 2 + topicNameBytes;
343             }
344 
345             int variablePartSize = variableHeaderBufferSize + payloadBufferSize;
346             int fixedHeaderBufferSize = 1 + getVariableLengthInt(variablePartSize);
347 
348             ByteBuf buf = ctx.alloc().buffer(fixedHeaderBufferSize + variablePartSize);
349             buf.writeByte(getFixedHeaderByte1(mqttFixedHeader));
350             writeVariableLengthInt(buf, variablePartSize);
351 
352             // Variable Header
353             int messageId = variableHeader.messageId();
354             buf.writeShort(messageId);
355             buf.writeBytes(propertiesBuf);
356 
357             // Payload
358             for (String topicName : payload.topics()) {
359                 writeEagerUTF8String(buf, topicName);
360             }
361 
362             return buf;
363         } finally {
364             propertiesBuf.release();
365         }
366     }
367 
368     private static ByteBuf encodeSubAckMessage(
369             ChannelHandlerContext ctx,
370             MqttSubAckMessage message) {
371         MqttVersion mqttVersion = getMqttVersion(ctx);
372         ByteBuf propertiesBuf = encodePropertiesIfNeeded(mqttVersion,
373                 ctx.alloc(),
374                 message.idAndPropertiesVariableHeader().properties());
375         try {
376             int variableHeaderBufferSize = 2 + propertiesBuf.readableBytes();
377             int payloadBufferSize = message.payload().grantedQoSLevels().size();
378             int variablePartSize = variableHeaderBufferSize + payloadBufferSize;
379             int fixedHeaderBufferSize = 1 + getVariableLengthInt(variablePartSize);
380             ByteBuf buf = ctx.alloc().buffer(fixedHeaderBufferSize + variablePartSize);
381             buf.writeByte(getFixedHeaderByte1(message.fixedHeader()));
382             writeVariableLengthInt(buf, variablePartSize);
383             buf.writeShort(message.variableHeader().messageId());
384             buf.writeBytes(propertiesBuf);
385             for (int code: message.payload().reasonCodes()) {
386                 buf.writeByte(code);
387             }
388 
389             return buf;
390         } finally {
391             propertiesBuf.release();
392         }
393     }
394 
395     private static ByteBuf encodeUnsubAckMessage(
396             ChannelHandlerContext ctx,
397             MqttUnsubAckMessage message) {
398         if (message.variableHeader() instanceof  MqttMessageIdAndPropertiesVariableHeader) {
399             MqttVersion mqttVersion = getMqttVersion(ctx);
400             // Reason Codes were introduced in MQTT 5.0 only. MQTT 3.1.1 (and 3.1) UNSUBACK packets
401             // have no payload, so reason codes must be suppressed for older protocol versions
402             // even when the caller populated them via MqttMessageBuilders.
403             final boolean writeReasonCodes = mqttVersion == MqttVersion.MQTT_5;
404             ByteBuf propertiesBuf = encodePropertiesIfNeeded(mqttVersion,
405                     ctx.alloc(),
406                     message.idAndPropertiesVariableHeader().properties());
407             try {
408                 int variableHeaderBufferSize = 2 + propertiesBuf.readableBytes();
409                 MqttUnsubAckPayload payload = message.payload();
410                 int payloadBufferSize = writeReasonCodes && payload != null
411                         ? payload.unsubscribeReasonCodes().size() : 0;
412                 int variablePartSize = variableHeaderBufferSize + payloadBufferSize;
413                 int fixedHeaderBufferSize = 1 + getVariableLengthInt(variablePartSize);
414                 ByteBuf buf = ctx.alloc().buffer(fixedHeaderBufferSize + variablePartSize);
415                 buf.writeByte(getFixedHeaderByte1(message.fixedHeader()));
416                 writeVariableLengthInt(buf, variablePartSize);
417                 buf.writeShort(message.variableHeader().messageId());
418                 buf.writeBytes(propertiesBuf);
419 
420                 if (writeReasonCodes && payload != null) {
421                     for (Short reasonCode : payload.unsubscribeReasonCodes()) {
422                         buf.writeByte(reasonCode);
423                     }
424                 }
425 
426                 return buf;
427             } finally {
428                 propertiesBuf.release();
429             }
430         } else {
431             return encodeMessageWithOnlySingleByteFixedHeaderAndMessageId(ctx.alloc(), message);
432         }
433     }
434 
435     private static ByteBuf encodePublishMessage(
436             ChannelHandlerContext ctx,
437             MqttPublishMessage message) {
438         MqttVersion mqttVersion = getMqttVersion(ctx);
439         MqttFixedHeader mqttFixedHeader = message.fixedHeader();
440         MqttPublishVariableHeader variableHeader = message.variableHeader();
441         ByteBuf payload = message.payload().duplicate();
442 
443         String topicName = variableHeader.topicName();
444         if (!isValidPublishTopicName(topicName)) {
445             throw new MqttIdentifierRejectedException("invalid topicName");
446         }
447         int topicNameBytes = utf8Bytes(topicName);
448 
449         ByteBuf propertiesBuf = encodePropertiesIfNeeded(mqttVersion,
450                 ctx.alloc(),
451                 message.variableHeader().properties());
452 
453         try {
454             int variableHeaderBufferSize = 2 + topicNameBytes +
455                     (mqttFixedHeader.qosLevel().value() > 0 ? 2 : 0) + propertiesBuf.readableBytes();
456             int payloadBufferSize = payload.readableBytes();
457             int variablePartSize = variableHeaderBufferSize + payloadBufferSize;
458             int fixedHeaderBufferSize = 1 + getVariableLengthInt(variablePartSize);
459 
460             ByteBuf buf = ctx.alloc().buffer(fixedHeaderBufferSize + variablePartSize);
461             buf.writeByte(getFixedHeaderByte1(mqttFixedHeader));
462             writeVariableLengthInt(buf, variablePartSize);
463             writeExactUTF8String(buf, topicName, topicNameBytes);
464             if (mqttFixedHeader.qosLevel().value() > 0) {
465                 buf.writeShort(variableHeader.packetId());
466             }
467             buf.writeBytes(propertiesBuf);
468             buf.writeBytes(payload);
469 
470             return buf;
471         } finally {
472             propertiesBuf.release();
473         }
474     }
475 
476     private static ByteBuf encodePubReplyMessage(ChannelHandlerContext ctx,
477                                           MqttMessage message) {
478         if (message.variableHeader() instanceof MqttPubReplyMessageVariableHeader) {
479             MqttFixedHeader mqttFixedHeader = message.fixedHeader();
480             MqttPubReplyMessageVariableHeader variableHeader =
481                     (MqttPubReplyMessageVariableHeader) message.variableHeader();
482             int msgId = variableHeader.messageId();
483 
484             final ByteBuf propertiesBuf;
485             final boolean includeReasonCode;
486             final int variableHeaderBufferSize;
487             final MqttVersion mqttVersion = getMqttVersion(ctx);
488             if (mqttVersion == MqttVersion.MQTT_5 &&
489                     (variableHeader.reasonCode() != MqttPubReplyMessageVariableHeader.REASON_CODE_OK ||
490                             !variableHeader.properties().isEmpty())) {
491                 propertiesBuf = encodeProperties(ctx.alloc(), variableHeader.properties());
492                 includeReasonCode = true;
493                 variableHeaderBufferSize = 3 + propertiesBuf.readableBytes();
494             } else {
495                 propertiesBuf = Unpooled.EMPTY_BUFFER;
496                 includeReasonCode = false;
497                 variableHeaderBufferSize = 2;
498             }
499 
500             try {
501                 final int fixedHeaderBufferSize = 1 + getVariableLengthInt(variableHeaderBufferSize);
502                 ByteBuf buf = ctx.alloc().buffer(fixedHeaderBufferSize + variableHeaderBufferSize);
503                 buf.writeByte(getFixedHeaderByte1(mqttFixedHeader));
504                 writeVariableLengthInt(buf, variableHeaderBufferSize);
505                 buf.writeShort(msgId);
506                 if (includeReasonCode) {
507                     buf.writeByte(variableHeader.reasonCode());
508                 }
509                 buf.writeBytes(propertiesBuf);
510 
511                 return buf;
512             } finally {
513                 propertiesBuf.release();
514             }
515         } else {
516             return encodeMessageWithOnlySingleByteFixedHeaderAndMessageId(ctx.alloc(), message);
517         }
518     }
519 
520     private static ByteBuf encodeMessageWithOnlySingleByteFixedHeaderAndMessageId(
521             ByteBufAllocator byteBufAllocator,
522             MqttMessage message) {
523         MqttFixedHeader mqttFixedHeader = message.fixedHeader();
524         MqttMessageIdVariableHeader variableHeader = (MqttMessageIdVariableHeader) message.variableHeader();
525         int msgId = variableHeader.messageId();
526 
527         int variableHeaderBufferSize = 2; // variable part only has a message id
528         int fixedHeaderBufferSize = 1 + getVariableLengthInt(variableHeaderBufferSize);
529         ByteBuf buf = byteBufAllocator.buffer(fixedHeaderBufferSize + variableHeaderBufferSize);
530         buf.writeByte(getFixedHeaderByte1(mqttFixedHeader));
531         writeVariableLengthInt(buf, variableHeaderBufferSize);
532         buf.writeShort(msgId);
533 
534         return buf;
535     }
536 
537     private static ByteBuf encodeReasonCodePlusPropertiesMessage(
538             ChannelHandlerContext ctx,
539             MqttMessage message) {
540         if (message.variableHeader() instanceof MqttReasonCodeAndPropertiesVariableHeader) {
541             MqttVersion mqttVersion = getMqttVersion(ctx);
542             MqttFixedHeader mqttFixedHeader = message.fixedHeader();
543             MqttReasonCodeAndPropertiesVariableHeader variableHeader =
544                     (MqttReasonCodeAndPropertiesVariableHeader) message.variableHeader();
545 
546             final ByteBuf propertiesBuf;
547             final boolean includeReasonCode;
548             final int variableHeaderBufferSize;
549             if (mqttVersion == MqttVersion.MQTT_5 &&
550                     (variableHeader.reasonCode() != MqttReasonCodeAndPropertiesVariableHeader.REASON_CODE_OK ||
551                             !variableHeader.properties().isEmpty())) {
552                 propertiesBuf = encodeProperties(ctx.alloc(), variableHeader.properties());
553                 includeReasonCode = true;
554                 variableHeaderBufferSize = 1 + propertiesBuf.readableBytes();
555             } else {
556                 propertiesBuf = Unpooled.EMPTY_BUFFER;
557                 includeReasonCode = false;
558                 variableHeaderBufferSize = 0;
559             }
560 
561             try {
562                 final int fixedHeaderBufferSize = 1 + getVariableLengthInt(variableHeaderBufferSize);
563                 ByteBuf buf = ctx.alloc().buffer(fixedHeaderBufferSize + variableHeaderBufferSize);
564                 buf.writeByte(getFixedHeaderByte1(mqttFixedHeader));
565                 writeVariableLengthInt(buf, variableHeaderBufferSize);
566                 if (includeReasonCode) {
567                     buf.writeByte(variableHeader.reasonCode());
568                 }
569                 buf.writeBytes(propertiesBuf);
570 
571                 return buf;
572             } finally {
573                 propertiesBuf.release();
574             }
575         } else {
576             return encodeMessageWithOnlySingleByteFixedHeader(ctx.alloc(), message);
577         }
578     }
579 
580     private static ByteBuf encodeMessageWithOnlySingleByteFixedHeader(
581             ByteBufAllocator byteBufAllocator,
582             MqttMessage message) {
583         MqttFixedHeader mqttFixedHeader = message.fixedHeader();
584         ByteBuf buf = byteBufAllocator.buffer(2);
585         buf.writeByte(getFixedHeaderByte1(mqttFixedHeader));
586         buf.writeByte(0);
587 
588         return buf;
589     }
590 
591     private static ByteBuf encodePropertiesIfNeeded(MqttVersion mqttVersion,
592                                              ByteBufAllocator byteBufAllocator,
593                                              MqttProperties mqttProperties) {
594         if (mqttVersion == MqttVersion.MQTT_5) {
595             return encodeProperties(byteBufAllocator, mqttProperties);
596         }
597         return Unpooled.EMPTY_BUFFER;
598     }
599 
600     private static ByteBuf encodeProperties(ByteBufAllocator byteBufAllocator,
601                                             MqttProperties mqttProperties) {
602         ByteBuf propertiesHeaderBuf = byteBufAllocator.buffer();
603         // encode also the Properties part
604         try {
605             ByteBuf propertiesBuf = byteBufAllocator.buffer();
606             try {
607                 for (MqttProperties.MqttProperty property : mqttProperties.listAll()) {
608                     MqttProperties.MqttPropertyType propertyType =
609                             MqttProperties.MqttPropertyType.valueOf(property.propertyId);
610                     switch (propertyType) {
611                         case PAYLOAD_FORMAT_INDICATOR:
612                         case REQUEST_PROBLEM_INFORMATION:
613                         case REQUEST_RESPONSE_INFORMATION:
614                         case MAXIMUM_QOS:
615                         case RETAIN_AVAILABLE:
616                         case WILDCARD_SUBSCRIPTION_AVAILABLE:
617                         case SUBSCRIPTION_IDENTIFIER_AVAILABLE:
618                         case SHARED_SUBSCRIPTION_AVAILABLE:
619                             writeVariableLengthInt(propertiesBuf, property.propertyId);
620                             final byte bytePropValue = ((MqttProperties.IntegerProperty) property).value.byteValue();
621                             propertiesBuf.writeByte(bytePropValue);
622                             break;
623                         case SERVER_KEEP_ALIVE:
624                         case RECEIVE_MAXIMUM:
625                         case TOPIC_ALIAS_MAXIMUM:
626                         case TOPIC_ALIAS:
627                             writeVariableLengthInt(propertiesBuf, property.propertyId);
628                             final short twoBytesInPropValue =
629                                     ((MqttProperties.IntegerProperty) property).value.shortValue();
630                             propertiesBuf.writeShort(twoBytesInPropValue);
631                             break;
632                         case PUBLICATION_EXPIRY_INTERVAL:
633                         case SESSION_EXPIRY_INTERVAL:
634                         case WILL_DELAY_INTERVAL:
635                         case MAXIMUM_PACKET_SIZE:
636                             writeVariableLengthInt(propertiesBuf, property.propertyId);
637                             final int fourBytesIntPropValue = ((MqttProperties.IntegerProperty) property).value;
638                             propertiesBuf.writeInt(fourBytesIntPropValue);
639                             break;
640                         case SUBSCRIPTION_IDENTIFIER:
641                             writeVariableLengthInt(propertiesBuf, property.propertyId);
642                             final int vbi = ((MqttProperties.IntegerProperty) property).value;
643                             writeVariableLengthInt(propertiesBuf, vbi);
644                             break;
645                         case CONTENT_TYPE:
646                         case RESPONSE_TOPIC:
647                         case ASSIGNED_CLIENT_IDENTIFIER:
648                         case AUTHENTICATION_METHOD:
649                         case RESPONSE_INFORMATION:
650                         case SERVER_REFERENCE:
651                         case REASON_STRING:
652                             writeVariableLengthInt(propertiesBuf, property.propertyId);
653                             writeEagerUTF8String(propertiesBuf, ((MqttProperties.StringProperty) property).value);
654                             break;
655                         case USER_PROPERTY:
656                             final List<MqttProperties.StringPair> pairs =
657                                     ((MqttProperties.UserProperties) property).value;
658                             for (MqttProperties.StringPair pair : pairs) {
659                                 writeVariableLengthInt(propertiesBuf, property.propertyId);
660                                 writeEagerUTF8String(propertiesBuf, pair.key);
661                                 writeEagerUTF8String(propertiesBuf, pair.value);
662                             }
663                             break;
664                         case CORRELATION_DATA:
665                         case AUTHENTICATION_DATA:
666                             writeVariableLengthInt(propertiesBuf, property.propertyId);
667                             final byte[] binaryPropValue = ((MqttProperties.BinaryProperty) property).value;
668                             propertiesBuf.writeShort(binaryPropValue.length);
669                             propertiesBuf.writeBytes(binaryPropValue, 0, binaryPropValue.length);
670                             break;
671                         default:
672                             //shouldn't reach here
673                             throw new EncoderException("Unknown property type: " + propertyType);
674                     }
675                 }
676                 writeVariableLengthInt(propertiesHeaderBuf, propertiesBuf.readableBytes());
677                 propertiesHeaderBuf.writeBytes(propertiesBuf);
678 
679                 return propertiesHeaderBuf;
680             } finally {
681                 propertiesBuf.release();
682             }
683         } catch (RuntimeException e) {
684             propertiesHeaderBuf.release();
685             throw e;
686         }
687     }
688 
689     private static int getFixedHeaderByte1(MqttFixedHeader header) {
690         int ret = 0;
691         ret |= header.messageType().value() << 4;
692         if (header.isDup()) {
693             ret |= 0x08;
694         }
695         ret |= header.qosLevel().value() << 1;
696         if (header.isRetain()) {
697             ret |= 0x01;
698         }
699         return ret;
700     }
701 
702     private static void writeVariableLengthInt(ByteBuf buf, int num) {
703         do {
704             int digit = num % 128;
705             num /= 128;
706             if (num > 0) {
707                 digit |= 0x80;
708             }
709             buf.writeByte(digit);
710         } while (num > 0);
711     }
712 
713     private static int nullableUtf8Bytes(String s) {
714         return s == null? 0 : utf8Bytes(s);
715     }
716 
717     private static int nullableMaxUtf8Bytes(String s) {
718         return s == null? 0 : utf8MaxBytes(s);
719     }
720 
721     private static void writeExactUTF8String(ByteBuf buf, String s, int utf8Length) {
722         buf.ensureWritable(utf8Length + 2);
723         buf.writeShort(utf8Length);
724         if (utf8Length > 0) {
725             final int writtenUtf8Length = reserveAndWriteUtf8(buf, s, utf8Length);
726             assert writtenUtf8Length == utf8Length;
727         }
728     }
729 
730     private static void writeEagerUTF8String(ByteBuf buf, String s) {
731         final int maxUtf8Length = nullableMaxUtf8Bytes(s);
732         buf.ensureWritable(maxUtf8Length + 2);
733         final int writerIndex = buf.writerIndex();
734         final int startUtf8String = writerIndex + 2;
735         buf.writerIndex(startUtf8String);
736         final int utf8Length = s != null? reserveAndWriteUtf8(buf, s, maxUtf8Length) : 0;
737         buf.setShort(writerIndex, utf8Length);
738     }
739 
740     private static int getVariableLengthInt(int num) {
741         int count = 0;
742         do {
743             num /= 128;
744             count++;
745         } while (num > 0);
746         return count;
747     }
748 
749 }