Browse Source

minor refactoring

pull/11450/head
YevhenBondarenko 2 years ago
parent
commit
b5fff2e44e
  1. 8
      application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/MqttClientSideRpcIntegrationTest.java
  2. 2
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java

8
application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/MqttClientSideRpcIntegrationTest.java

@ -49,8 +49,8 @@ public class MqttClientSideRpcIntegrationTest extends AbstractMqttIntegrationTes
@Value("${transport.mqtt.netty.max_payload_size}")
private Integer maxPayloadSize;
@Value("${transport.mqtt.msg_queue_size_per_device_limit:100}")
private int maxQueueSize;
@Value("${transport.mqtt.msg_queue_size_per_device_limit}")
private int maxInflightMessages;
@Before
public void beforeTest() throws Exception {
@ -110,7 +110,7 @@ public class MqttClientSideRpcIntegrationTest extends AbstractMqttIntegrationTes
assertEquals(response.get("deviceTelemetryMsgRateLimit"), profileConfiguration.getTransportDeviceTelemetryMsgRateLimit());
assertEquals(response.get("deviceTelemetryDataPointsRateLimit"), profileConfiguration.getTransportDeviceTelemetryDataPointsRateLimit());
assertEquals(response.get("maxPayloadSize"), maxPayloadSize);
assertEquals(response.get("maxQueueSize"), maxQueueSize);
assertEquals(response.get("maxInflightMessages"), maxInflightMessages);
assertEquals(response.get("payloadType"), TransportPayloadType.JSON.name());
client.disconnect();
@ -160,7 +160,7 @@ public class MqttClientSideRpcIntegrationTest extends AbstractMqttIntegrationTes
assertEquals(response.get("gatewayDeviceTelemetryMsgRateLimit"), profileConfiguration.getTransportGatewayDeviceTelemetryMsgRateLimit());
assertEquals(response.get("gatewayDeviceTelemetryDataPointsRateLimit"), profileConfiguration.getTransportGatewayDeviceTelemetryDataPointsRateLimit());
assertEquals(response.get("maxPayloadSize"), maxPayloadSize);
assertEquals(response.get("maxQueueSize"), maxQueueSize);
assertEquals(response.get("maxInflightMessages"), maxInflightMessages);
assertEquals(response.get("payloadType"), TransportPayloadType.JSON.name());
client.disconnect();

2
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java

@ -1362,7 +1362,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
serviceConfiguration.put("deviceTelemetryDataPointsRateLimit", profile.getTransportDeviceTelemetryDataPointsRateLimit());
}
serviceConfiguration.put("maxPayloadSize", context.getMaxPayloadSize());
serviceConfiguration.put("maxQueueSize", context.getMessageQueueSizePerDeviceLimit());
serviceConfiguration.put("maxInflightMessages", context.getMessageQueueSizePerDeviceLimit());
serviceConfiguration.put("payloadType", deviceSessionCtx.getPayloadType());
ack(ctx, msgId, MqttReasonCodes.PubAck.SUCCESS);

Loading…
Cancel
Save