Browse Source

Suppress spaming logs and fix connection statistics in case of rate limits by ip

pull/6174/head
Andrii Shvaika 5 years ago
parent
commit
40c2fe3fd2
  1. 2
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/JsonCoapAdaptor.java
  2. 16
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java
  3. 14
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java
  4. 11
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/limits/ProxyIpFilter.java

2
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 { private String validatePayload(UUID sessionId, Request inbound, boolean isEmptyPayloadAllowed) throws AdaptorException {
String payload = inbound.getPayloadString(); String payload = inbound.getPayloadString();
if (payload == null) { if (payload == null) {
log.warn("[{}] Payload is empty!", sessionId); log.debug("[{}] Payload is empty!", sessionId);
if (!isEmptyPayloadAllowed) { if (!isEmptyPayloadAllowed) {
throw new AdaptorException(new IllegalArgumentException("Payload is empty!")); throw new AdaptorException(new IllegalArgumentException("Payload is empty!"));
} }

16
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java

@ -61,7 +61,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor {
try { try {
return JsonConverter.convertToTelemetryProto(new JsonParser().parse(payload)); return JsonConverter.convertToTelemetryProto(new JsonParser().parse(payload));
} catch (IllegalStateException | JsonSyntaxException ex) { } 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); throw new AdaptorException(ex);
} }
} }
@ -72,7 +72,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor {
try { try {
return JsonConverter.convertToAttributesProto(new JsonParser().parse(payload)); return JsonConverter.convertToAttributesProto(new JsonParser().parse(payload));
} catch (IllegalStateException | JsonSyntaxException ex) { } 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); throw new AdaptorException(ex);
} }
} }
@ -83,7 +83,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor {
try { try {
return JsonConverter.convertToClaimDeviceProto(ctx.getDeviceId(), payload); return JsonConverter.convertToClaimDeviceProto(ctx.getDeviceId(), payload);
} catch (IllegalStateException | JsonSyntaxException ex) { } 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); throw new AdaptorException(ex);
} }
} }
@ -164,7 +164,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor {
try { try {
return new JsonParser().parse(payload); return new JsonParser().parse(payload);
} catch (JsonSyntaxException ex) { } catch (JsonSyntaxException ex) {
log.warn("Payload is in incorrect format: {}", payload); log.debug("Payload is in incorrect format: {}", payload);
throw new AdaptorException(ex); throw new AdaptorException(ex);
} }
} }
@ -186,7 +186,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor {
} }
return result.build(); return result.build();
} catch (RuntimeException e) { } 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); throw new AdaptorException(e);
} }
} }
@ -198,7 +198,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor {
String payload = inbound.payload().toString(UTF8); String payload = inbound.payload().toString(UTF8);
return TransportProtos.ToDeviceRpcResponseMsg.newBuilder().setRequestId(requestId).setPayload(payload).build(); return TransportProtos.ToDeviceRpcResponseMsg.newBuilder().setRequestId(requestId).setPayload(payload).build();
} catch (RuntimeException e) { } catch (RuntimeException e) {
log.warn("Failed to decode rpc response", e); log.debug("Failed to decode rpc response", e);
throw new AdaptorException(e); throw new AdaptorException(e);
} }
} }
@ -210,7 +210,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor {
int requestId = getRequestId(topicName, topicBase); int requestId = getRequestId(topicName, topicBase);
return JsonConverter.convertToServerRpcRequest(new JsonParser().parse(payload), requestId); return JsonConverter.convertToServerRpcRequest(new JsonParser().parse(payload), requestId);
} catch (IllegalStateException | JsonSyntaxException ex) { } 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); 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 { private static String validatePayload(UUID sessionId, ByteBuf payloadData, boolean isEmptyPayloadAllowed) throws AdaptorException {
String payload = payloadData.toString(UTF8); String payload = payloadData.toString(UTF8);
if (payload == null) { if (payload == null) {
log.warn("[{}] Payload is empty!", sessionId); log.debug("[{}] Payload is empty!", sessionId);
if (!isEmptyPayloadAllowed) { if (!isEmptyPayloadAllowed) {
throw new AdaptorException(new IllegalArgumentException("Payload is empty!")); throw new AdaptorException(new IllegalArgumentException("Payload is empty!"));
} }

14
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java

@ -52,7 +52,7 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor {
try { try {
return JsonConverter.convertToTelemetryProto(new JsonParser().parse(ProtoConverter.dynamicMsgToJson(bytes, telemetryDynamicMsgDescriptor))); return JsonConverter.convertToTelemetryProto(new JsonParser().parse(ProtoConverter.dynamicMsgToJson(bytes, telemetryDynamicMsgDescriptor)));
} catch (Exception e) { } 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); throw new AdaptorException(e);
} }
} }
@ -65,7 +65,7 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor {
try { try {
return JsonConverter.convertToAttributesProto(new JsonParser().parse(ProtoConverter.dynamicMsgToJson(bytes, attributesDynamicMessageDescriptor))); return JsonConverter.convertToAttributesProto(new JsonParser().parse(ProtoConverter.dynamicMsgToJson(bytes, attributesDynamicMessageDescriptor)));
} catch (Exception e) { } 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); throw new AdaptorException(e);
} }
} }
@ -76,7 +76,7 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor {
try { try {
return ProtoConverter.convertToClaimDeviceProto(ctx.getDeviceId(), bytes); return ProtoConverter.convertToClaimDeviceProto(ctx.getDeviceId(), bytes);
} catch (InvalidProtocolBufferException e) { } 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); throw new AdaptorException(e);
} }
} }
@ -89,7 +89,7 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor {
int requestId = getRequestId(topicName, topicBase); int requestId = getRequestId(topicName, topicBase);
return ProtoConverter.convertToGetAttributeRequestMessage(bytes, requestId); return ProtoConverter.convertToGetAttributeRequestMessage(bytes, requestId);
} catch (InvalidProtocolBufferException e) { } 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); throw new AdaptorException(e);
} }
} }
@ -105,7 +105,7 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor {
JsonElement response = new JsonParser().parse(ProtoConverter.dynamicMsgToJson(bytes, rpcResponseDynamicMessageDescriptor)); JsonElement response = new JsonParser().parse(ProtoConverter.dynamicMsgToJson(bytes, rpcResponseDynamicMessageDescriptor));
return TransportProtos.ToDeviceRpcResponseMsg.newBuilder().setRequestId(requestId).setPayload(response.toString()).build(); return TransportProtos.ToDeviceRpcResponseMsg.newBuilder().setRequestId(requestId).setPayload(response.toString()).build();
} catch (Exception e) { } catch (Exception e) {
log.warn("Failed to decode rpc response", e); log.debug("Failed to decode rpc response", e);
throw new AdaptorException(e); throw new AdaptorException(e);
} }
} }
@ -118,7 +118,7 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor {
int requestId = getRequestId(topicName, topicBase); int requestId = getRequestId(topicName, topicBase);
return ProtoConverter.convertToServerRpcRequest(bytes, requestId); return ProtoConverter.convertToServerRpcRequest(bytes, requestId);
} catch (InvalidProtocolBufferException e) { } 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); throw new AdaptorException(e);
} }
} }
@ -129,7 +129,7 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor {
try { try {
return ProtoConverter.convertToProvisionRequestMsg(bytes); return ProtoConverter.convertToProvisionRequestMsg(bytes);
} catch (InvalidProtocolBufferException ex) { } catch (InvalidProtocolBufferException ex) {
log.warn("Failed to decode provision request", ex); log.debug("Failed to decode provision request", ex);
throw new AdaptorException(ex); throw new AdaptorException(ex);
} }
} }

11
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; package org.thingsboard.server.transport.mqtt.limits;
import io.netty.channel.ChannelHandler;
import io.netty.channel.ChannelHandlerContext; import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelInboundHandlerAdapter; import io.netty.channel.ChannelInboundHandlerAdapter;
import io.netty.handler.codec.haproxy.HAProxyMessage; import io.netty.handler.codec.haproxy.HAProxyMessage;
@ -60,7 +61,15 @@ public class ProxyIpFilter extends ChannelInboundHandlerAdapter {
private void closeChannel(ChannelHandlerContext ctx) { private void closeChannel(ChannelHandlerContext ctx) {
while (ctx.pipeline().last() != this) { 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.pipeline().remove(this);
ctx.close(); ctx.close();

Loading…
Cancel
Save