diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/JsonCoapAdaptor.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/JsonCoapAdaptor.java index 45e5009a7b..3d4a4ee053 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/JsonCoapAdaptor.java +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/JsonCoapAdaptor.java @@ -153,7 +153,7 @@ public class JsonCoapAdaptor implements CoapTransportAdaptor { private String validatePayload(UUID sessionId, Request inbound, boolean isEmptyPayloadAllowed) throws AdaptorException { String payload = inbound.getPayloadString(); if (payload == null) { - log.warn("[{}] Payload is empty!", sessionId); + log.debug("[{}] Payload is empty!", sessionId); if (!isEmptyPayloadAllowed) { throw new AdaptorException(new IllegalArgumentException("Payload is empty!")); } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java index 73c0e9a32f..75c76c0cdb 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java @@ -61,7 +61,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { try { return JsonConverter.convertToTelemetryProto(new JsonParser().parse(payload)); } catch (IllegalStateException | JsonSyntaxException ex) { - log.warn("Failed to decode post telemetry request", ex); + log.debug("Failed to decode post telemetry request", ex); throw new AdaptorException(ex); } } @@ -72,7 +72,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { try { return JsonConverter.convertToAttributesProto(new JsonParser().parse(payload)); } catch (IllegalStateException | JsonSyntaxException ex) { - log.warn("Failed to decode post attributes request", ex); + log.debug("Failed to decode post attributes request", ex); throw new AdaptorException(ex); } } @@ -83,7 +83,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { try { return JsonConverter.convertToClaimDeviceProto(ctx.getDeviceId(), payload); } catch (IllegalStateException | JsonSyntaxException ex) { - log.warn("Failed to decode claim device request", ex); + log.debug("Failed to decode claim device request", ex); throw new AdaptorException(ex); } } @@ -164,7 +164,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { try { return new JsonParser().parse(payload); } catch (JsonSyntaxException ex) { - log.warn("Payload is in incorrect format: {}", payload); + log.debug("Payload is in incorrect format: {}", payload); throw new AdaptorException(ex); } } @@ -186,7 +186,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { } return result.build(); } catch (RuntimeException e) { - log.warn("Failed to decode get attributes request", e); + log.debug("Failed to decode get attributes request", e); throw new AdaptorException(e); } } @@ -198,7 +198,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { String payload = inbound.payload().toString(UTF8); return TransportProtos.ToDeviceRpcResponseMsg.newBuilder().setRequestId(requestId).setPayload(payload).build(); } catch (RuntimeException e) { - log.warn("Failed to decode rpc response", e); + log.debug("Failed to decode rpc response", e); throw new AdaptorException(e); } } @@ -210,7 +210,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { int requestId = getRequestId(topicName, topicBase); return JsonConverter.convertToServerRpcRequest(new JsonParser().parse(payload), requestId); } catch (IllegalStateException | JsonSyntaxException ex) { - log.warn("Failed to decode to server rpc request", ex); + log.debug("Failed to decode to server rpc request", ex); throw new AdaptorException(ex); } } @@ -259,7 +259,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { private static String validatePayload(UUID sessionId, ByteBuf payloadData, boolean isEmptyPayloadAllowed) throws AdaptorException { String payload = payloadData.toString(UTF8); if (payload == null) { - log.warn("[{}] Payload is empty!", sessionId); + log.debug("[{}] Payload is empty!", sessionId); if (!isEmptyPayloadAllowed) { throw new AdaptorException(new IllegalArgumentException("Payload is empty!")); } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java index 1f0aef1c64..175bf17c02 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java @@ -52,7 +52,7 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor { try { return JsonConverter.convertToTelemetryProto(new JsonParser().parse(ProtoConverter.dynamicMsgToJson(bytes, telemetryDynamicMsgDescriptor))); } catch (Exception e) { - log.warn("Failed to decode post telemetry request", e); + log.debug("Failed to decode post telemetry request", e); throw new AdaptorException(e); } } @@ -65,7 +65,7 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor { try { return JsonConverter.convertToAttributesProto(new JsonParser().parse(ProtoConverter.dynamicMsgToJson(bytes, attributesDynamicMessageDescriptor))); } catch (Exception e) { - log.warn("Failed to decode post attributes request", e); + log.debug("Failed to decode post attributes request", e); throw new AdaptorException(e); } } @@ -76,7 +76,7 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor { try { return ProtoConverter.convertToClaimDeviceProto(ctx.getDeviceId(), bytes); } catch (InvalidProtocolBufferException e) { - log.warn("Failed to decode claim device request", e); + log.debug("Failed to decode claim device request", e); throw new AdaptorException(e); } } @@ -89,7 +89,7 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor { int requestId = getRequestId(topicName, topicBase); return ProtoConverter.convertToGetAttributeRequestMessage(bytes, requestId); } catch (InvalidProtocolBufferException e) { - log.warn("Failed to decode get attributes request", e); + log.debug("Failed to decode get attributes request", e); throw new AdaptorException(e); } } @@ -105,7 +105,7 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor { JsonElement response = new JsonParser().parse(ProtoConverter.dynamicMsgToJson(bytes, rpcResponseDynamicMessageDescriptor)); return TransportProtos.ToDeviceRpcResponseMsg.newBuilder().setRequestId(requestId).setPayload(response.toString()).build(); } catch (Exception e) { - log.warn("Failed to decode rpc response", e); + log.debug("Failed to decode rpc response", e); throw new AdaptorException(e); } } @@ -118,7 +118,7 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor { int requestId = getRequestId(topicName, topicBase); return ProtoConverter.convertToServerRpcRequest(bytes, requestId); } catch (InvalidProtocolBufferException e) { - log.warn("Failed to decode to server rpc request", e); + log.debug("Failed to decode to server rpc request", e); throw new AdaptorException(e); } } @@ -129,7 +129,7 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor { try { return ProtoConverter.convertToProvisionRequestMsg(bytes); } catch (InvalidProtocolBufferException ex) { - log.warn("Failed to decode provision request", ex); + log.debug("Failed to decode provision request", ex); throw new AdaptorException(ex); } } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/limits/ProxyIpFilter.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/limits/ProxyIpFilter.java index c7a70c6d44..c6f605c9dd 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/limits/ProxyIpFilter.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/limits/ProxyIpFilter.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.transport.mqtt.limits; +import io.netty.channel.ChannelHandler; import io.netty.channel.ChannelHandlerContext; import io.netty.channel.ChannelInboundHandlerAdapter; import io.netty.handler.codec.haproxy.HAProxyMessage; @@ -60,7 +61,15 @@ public class ProxyIpFilter extends ChannelInboundHandlerAdapter { private void closeChannel(ChannelHandlerContext ctx) { while (ctx.pipeline().last() != this) { - ctx.pipeline().removeLast(); + ChannelHandler handler = ctx.pipeline().removeLast(); + if (handler instanceof ChannelInboundHandlerAdapter) { + try { + ((ChannelInboundHandlerAdapter) handler).channelUnregistered(ctx); + } catch (Exception e) { + log.error("Failed to unregister channel: [{}]", ctx, e); + } + } + } ctx.pipeline().remove(this); ctx.close();