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
38
39
40
41
42
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
58
59
60
61
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
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
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
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
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
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
293 int messageId = variableHeader.messageId();
294 buf.writeShort(messageId);
295 buf.writeBytes(propertiesBuf);
296
297
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
353 int messageId = variableHeader.messageId();
354 buf.writeShort(messageId);
355 buf.writeBytes(propertiesBuf);
356
357
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
401
402
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;
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
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
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 }