|
|
|
@ -41,6 +41,7 @@ import java.util.function.BiConsumer; |
|
|
|
import java.util.function.Function; |
|
|
|
|
|
|
|
import static org.thingsboard.server.common.transport.limits.TransportLimitsType.DEVICE_LIMITS; |
|
|
|
import static org.thingsboard.server.common.transport.limits.TransportLimitsType.GATEWAY_DEVICE_LIMITS; |
|
|
|
import static org.thingsboard.server.common.transport.limits.TransportLimitsType.GATEWAY_LIMITS; |
|
|
|
import static org.thingsboard.server.common.transport.limits.TransportLimitsType.TENANT_LIMITS; |
|
|
|
|
|
|
|
@ -53,9 +54,11 @@ public class DefaultTransportRateLimitService implements TransportRateLimitServi |
|
|
|
private final ConcurrentMap<TenantId, Boolean> tenantAllowed = new ConcurrentHashMap<>(); |
|
|
|
private final ConcurrentMap<TenantId, Set<DeviceId>> tenantDevices = new ConcurrentHashMap<>(); |
|
|
|
private final ConcurrentMap<TenantId, Set<DeviceId>> tenantGateways = new ConcurrentHashMap<>(); |
|
|
|
private final ConcurrentMap<TenantId, Set<DeviceId>> tenantGatewayDevices = new ConcurrentHashMap<>(); |
|
|
|
private final ConcurrentMap<TenantId, EntityTransportRateLimits> perTenantLimits = new ConcurrentHashMap<>(); |
|
|
|
private final ConcurrentMap<DeviceId, EntityTransportRateLimits> perDeviceLimits = new ConcurrentHashMap<>(); |
|
|
|
private final ConcurrentMap<DeviceId, EntityTransportRateLimits> perGatewayLimits = new ConcurrentHashMap<>(); |
|
|
|
private final ConcurrentMap<DeviceId, EntityTransportRateLimits> perGatewayDeviceLimits = new ConcurrentHashMap<>(); |
|
|
|
private final Map<InetAddress, InetAddressRateLimitStats> ipMap = new ConcurrentHashMap<>(); |
|
|
|
|
|
|
|
private final TransportTenantProfileCache tenantProfileCache; |
|
|
|
@ -72,21 +75,23 @@ public class DefaultTransportRateLimitService implements TransportRateLimitServi |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public TbPair<EntityType, Boolean> checkLimits(TenantId tenantId, DeviceId gatewayId, DeviceId deviceId, int dataPoints) { |
|
|
|
public TbPair<EntityType, Boolean> checkLimits(TenantId tenantId, DeviceId gatewayId, DeviceId deviceId, int dataPoints, boolean isGateway) { |
|
|
|
if (!tenantAllowed.getOrDefault(tenantId, Boolean.TRUE)) { |
|
|
|
return TbPair.of(EntityType.API_USAGE_STATE, false); |
|
|
|
} |
|
|
|
if (!checkEntityRateLimit(dataPoints, getTenantRateLimits(tenantId))) { |
|
|
|
return TbPair.of(EntityType.TENANT, false); |
|
|
|
} |
|
|
|
|
|
|
|
if (isGateway && !checkEntityRateLimit(dataPoints, getGatewayDeviceRateLimits(tenantId, deviceId))) { |
|
|
|
return TbPair.of(EntityType.DEVICE, true); |
|
|
|
} |
|
|
|
if (gatewayId != null && !checkEntityRateLimit(dataPoints, getGatewayRateLimits(tenantId, gatewayId))) { |
|
|
|
return TbPair.of(EntityType.DEVICE, true); |
|
|
|
} |
|
|
|
|
|
|
|
if (deviceId != null && !checkEntityRateLimit(dataPoints, getDeviceRateLimits(tenantId, deviceId))) { |
|
|
|
if (!isGateway && deviceId != null && !checkEntityRateLimit(dataPoints, getDeviceRateLimits(tenantId, deviceId))) { |
|
|
|
return TbPair.of(EntityType.DEVICE, false); |
|
|
|
} |
|
|
|
|
|
|
|
return null; |
|
|
|
} |
|
|
|
|
|
|
|
@ -104,8 +109,9 @@ public class DefaultTransportRateLimitService implements TransportRateLimitServi |
|
|
|
EntityTransportRateLimits tenantRateLimitPrototype = createRateLimits(update.getProfile(), TENANT_LIMITS); |
|
|
|
EntityTransportRateLimits deviceRateLimitPrototype = createRateLimits(update.getProfile(), DEVICE_LIMITS); |
|
|
|
EntityTransportRateLimits gatewayRateLimitPrototype = createRateLimits(update.getProfile(), GATEWAY_LIMITS); |
|
|
|
EntityTransportRateLimits gatewayDeviceRateLimitPrototype = createRateLimits(update.getProfile(), GATEWAY_DEVICE_LIMITS); |
|
|
|
for (TenantId tenantId : update.getAffectedTenants()) { |
|
|
|
update(tenantId, tenantRateLimitPrototype, deviceRateLimitPrototype, gatewayRateLimitPrototype); |
|
|
|
update(tenantId, tenantRateLimitPrototype, deviceRateLimitPrototype, gatewayRateLimitPrototype, gatewayDeviceRateLimitPrototype); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@ -114,26 +120,34 @@ public class DefaultTransportRateLimitService implements TransportRateLimitServi |
|
|
|
EntityTransportRateLimits tenantRateLimitPrototype = createRateLimits(tenantProfileCache.get(tenantId), TENANT_LIMITS); |
|
|
|
EntityTransportRateLimits deviceRateLimitPrototype = createRateLimits(tenantProfileCache.get(tenantId), DEVICE_LIMITS); |
|
|
|
EntityTransportRateLimits gatewayRateLimitPrototype = createRateLimits(tenantProfileCache.get(tenantId), GATEWAY_LIMITS); |
|
|
|
update(tenantId, tenantRateLimitPrototype, deviceRateLimitPrototype, gatewayRateLimitPrototype); |
|
|
|
EntityTransportRateLimits gatewayDeviceRateLimitPrototype = createRateLimits(tenantProfileCache.get(tenantId), GATEWAY_DEVICE_LIMITS); |
|
|
|
update(tenantId, tenantRateLimitPrototype, deviceRateLimitPrototype, gatewayRateLimitPrototype, gatewayDeviceRateLimitPrototype); |
|
|
|
} |
|
|
|
|
|
|
|
private void update(TenantId tenantId, EntityTransportRateLimits tenantRateLimitPrototype, |
|
|
|
EntityTransportRateLimits deviceRateLimitPrototype, EntityTransportRateLimits gatewayRateLimitPrototype) { |
|
|
|
private void update(TenantId tenantId, EntityTransportRateLimits tenantRateLimitPrototype, EntityTransportRateLimits deviceRateLimitPrototype, |
|
|
|
EntityTransportRateLimits gatewayRateLimitPrototype, EntityTransportRateLimits gatewayDeviceRateLimitPrototype) { |
|
|
|
mergeLimits(tenantId, tenantRateLimitPrototype, perTenantLimits::get, perTenantLimits::put); |
|
|
|
getTenantDevices(tenantId).forEach(deviceId -> mergeLimits(deviceId, deviceRateLimitPrototype, perDeviceLimits::get, perDeviceLimits::put)); |
|
|
|
getTenantGateways(tenantId).forEach(deviceId -> mergeLimits(deviceId, gatewayRateLimitPrototype, perGatewayLimits::get, perGatewayLimits::put)); |
|
|
|
getTenantGateways(tenantId).forEach(gatewayId -> mergeLimits(gatewayId, gatewayRateLimitPrototype, perGatewayLimits::get, perGatewayLimits::put)); |
|
|
|
getTenantGatewayDevices(tenantId).forEach(gatewayId -> mergeLimits(gatewayId, gatewayDeviceRateLimitPrototype, perGatewayDeviceLimits::get, perGatewayDeviceLimits::put)); |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void remove(TenantId tenantId) { |
|
|
|
perTenantLimits.remove(tenantId); |
|
|
|
tenantDevices.remove(tenantId); |
|
|
|
tenantGateways.remove(tenantId); |
|
|
|
tenantGatewayDevices.remove(tenantId); |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void remove(DeviceId deviceId) { |
|
|
|
perDeviceLimits.remove(deviceId); |
|
|
|
perGatewayLimits.remove(deviceId); |
|
|
|
perGatewayDeviceLimits.remove(deviceId); |
|
|
|
tenantDevices.values().forEach(set -> set.remove(deviceId)); |
|
|
|
tenantGateways.values().forEach(set -> set.remove(deviceId)); |
|
|
|
tenantGatewayDevices.values().forEach(set -> set.remove(deviceId)); |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
@ -273,6 +287,11 @@ public class DefaultTransportRateLimitService implements TransportRateLimitServi |
|
|
|
telemetryMsgRateLimit = newLimit(profile.getTransportGatewayTelemetryMsgRateLimit()); |
|
|
|
telemetryDpRateLimit = newLimit(profile.getTransportGatewayTelemetryDataPointsRateLimit()); |
|
|
|
} |
|
|
|
case GATEWAY_DEVICE_LIMITS -> { |
|
|
|
regularMsgRateLimit = newLimit(profile.getTransportGatewayDeviceMsgRateLimit()); |
|
|
|
telemetryMsgRateLimit = newLimit(profile.getTransportGatewayDeviceTelemetryMsgRateLimit()); |
|
|
|
telemetryDpRateLimit = newLimit(profile.getTransportGatewayDeviceTelemetryDataPointsRateLimit()); |
|
|
|
} |
|
|
|
default -> throw new IllegalStateException("Unknown limits type: " + limitsType); |
|
|
|
} |
|
|
|
|
|
|
|
@ -296,10 +315,18 @@ public class DefaultTransportRateLimitService implements TransportRateLimitServi |
|
|
|
}); |
|
|
|
} |
|
|
|
|
|
|
|
private EntityTransportRateLimits getGatewayRateLimits(TenantId tenantId, DeviceId deviceId) { |
|
|
|
return perGatewayLimits.computeIfAbsent(deviceId, k -> { |
|
|
|
private EntityTransportRateLimits getGatewayRateLimits(TenantId tenantId, DeviceId gatewayId) { |
|
|
|
return perGatewayLimits.computeIfAbsent(gatewayId, k -> { |
|
|
|
EntityTransportRateLimits limits = createRateLimits(tenantProfileCache.get(tenantId), GATEWAY_LIMITS); |
|
|
|
getTenantGateways(tenantId).add(deviceId); |
|
|
|
getTenantGateways(tenantId).add(gatewayId); |
|
|
|
return limits; |
|
|
|
}); |
|
|
|
} |
|
|
|
|
|
|
|
private EntityTransportRateLimits getGatewayDeviceRateLimits(TenantId tenantId, DeviceId gatewayId) { |
|
|
|
return perGatewayDeviceLimits.computeIfAbsent(gatewayId, k -> { |
|
|
|
EntityTransportRateLimits limits = createRateLimits(tenantProfileCache.get(tenantId), GATEWAY_DEVICE_LIMITS); |
|
|
|
getTenantGatewayDevices(tenantId).add(gatewayId); |
|
|
|
return limits; |
|
|
|
}); |
|
|
|
} |
|
|
|
@ -312,4 +339,8 @@ public class DefaultTransportRateLimitService implements TransportRateLimitServi |
|
|
|
return tenantGateways.computeIfAbsent(tenantId, id -> ConcurrentHashMap.newKeySet()); |
|
|
|
} |
|
|
|
|
|
|
|
private Set<DeviceId> getTenantGatewayDevices(TenantId tenantId) { |
|
|
|
return tenantGatewayDevices.computeIfAbsent(tenantId, id -> ConcurrentHashMap.newKeySet()); |
|
|
|
} |
|
|
|
|
|
|
|
} |
|
|
|
|