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