From e9f5fd27066f2bfeaf1316d70d83d59c7e52a0cf Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Tue, 25 May 2021 17:02:21 +0300 Subject: [PATCH] added ability to use redis cluster for lwm2m redis stores. --- application/pom.xml | 4 + .../data/device/profile/MqttTopics.java | 2 +- common/transport/lwm2m/pom.xml | 4 + .../server/DefaultLwM2mTransportService.java | 5 +- .../store/TbLwM2mRedisRegistrationStore.java | 373 ++++++++---------- .../store/TbLwM2mRedisSecurityStore.java | 64 +-- .../server/store/TbLwM2mStoreFactory.java | 6 +- pom.xml | 5 + .../src/main/resources/tb-lwm2m-transport.yml | 3 +- 9 files changed, 230 insertions(+), 236 deletions(-) diff --git a/application/pom.xml b/application/pom.xml index e7cedeb540..303f206601 100644 --- a/application/pom.xml +++ b/application/pom.xml @@ -145,6 +145,10 @@ ${project.version} runtime + + org.springframework.integration + spring-integration-redis + org.springframework.boot spring-boot-starter-security diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttTopics.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttTopics.java index 6373d41803..7b634bffdc 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttTopics.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttTopics.java @@ -82,7 +82,7 @@ public class MqttTopics { public static final String DEVICE_FIRMWARE_REQUEST_TOPIC_PATTERN = BASE_DEVICE_API_TOPIC_V2 + FIRMWARE + REQUEST + "/" + REQUEST_ID_PATTERN + CHUNK + CHUNK_PATTERN; public static final String DEVICE_FIRMWARE_RESPONSES_TOPIC = BASE_DEVICE_API_TOPIC_V2 + FIRMWARE + RESPONSE + "/" + SUB_TOPIC + CHUNK + SUB_TOPIC; public static final String DEVICE_FIRMWARE_ERROR_TOPIC = BASE_DEVICE_API_TOPIC_V2 + FIRMWARE + ERROR; - public static final String DEVICE_FIRMWARE_RESPONSES_TOPIC_FORMAT = BASE_DEVICE_API_TOPIC_V2 + "%s" + RESPONSE + "/"+ "%s" + CHUNK + "%d"; + public static final String DEVICE_FIRMWARE_RESPONSES_TOPIC_FORMAT = BASE_DEVICE_API_TOPIC_V2 + "/%s" + RESPONSE + "/"+ "%s" + CHUNK + "%d"; public static final String DEVICE_SOFTWARE_REQUEST_TOPIC_PATTERN = BASE_DEVICE_API_TOPIC_V2 + SOFTWARE + REQUEST + "/" + REQUEST_ID_PATTERN + CHUNK + CHUNK_PATTERN; public static final String DEVICE_SOFTWARE_RESPONSES_TOPIC = BASE_DEVICE_API_TOPIC_V2 + SOFTWARE + RESPONSE + "/" + SUB_TOPIC + CHUNK + SUB_TOPIC; diff --git a/common/transport/lwm2m/pom.xml b/common/transport/lwm2m/pom.xml index aae103fc4f..1b639af89d 100644 --- a/common/transport/lwm2m/pom.xml +++ b/common/transport/lwm2m/pom.xml @@ -52,6 +52,10 @@ org.springframework spring-context + + org.springframework.integration + spring-integration-redis + org.slf4j slf4j-api diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2mTransportService.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2mTransportService.java index e23d55f3a1..c058712234 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2mTransportService.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2mTransportService.java @@ -24,7 +24,6 @@ import org.eclipse.leshan.core.util.Hex; import org.eclipse.leshan.server.californium.LeshanServer; import org.eclipse.leshan.server.californium.LeshanServerBuilder; import org.eclipse.leshan.server.californium.registration.CaliforniumRegistrationStore; -import org.eclipse.leshan.server.californium.registration.InMemoryRegistrationStore; import org.eclipse.leshan.server.model.LwM2mModelProvider; import org.eclipse.leshan.server.security.DefaultAuthorizer; import org.eclipse.leshan.server.security.EditableSecurityStore; @@ -127,8 +126,8 @@ public class DefaultLwM2mTransportService implements LwM2MTransportService { builder.setEncoder(new DefaultLwM2mNodeEncoder(LwM2mValueConverterImpl.getInstance())); /* InMemoryRegistrationStore(ScheduledExecutorService schedExecutor, long cleanPeriodInSec) */ - InMemoryRegistrationStore registrationStore = new InMemoryRegistrationStore(this.registrationStoreExecutor, this.config.getCleanPeriodInSec()); - builder.setRegistrationStore(registrationStore); +//// InMemoryRegistrationStore registrationStore = new InMemoryRegistrationStore(this.registrationStoreExecutor, this.config.getCleanPeriodInSec()); +// builder.setRegistrationStore(registrationStore); /* Create CoAP Config */ builder.setCoapConfig(getCoapConfig(config.getPort(), config.getSecurePort())); diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mRedisRegistrationStore.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mRedisRegistrationStore.java index d947e22133..1de2c1a4aa 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mRedisRegistrationStore.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mRedisRegistrationStore.java @@ -26,9 +26,7 @@ import org.eclipse.leshan.core.util.NamedThreadFactory; import org.eclipse.leshan.core.util.Validate; import org.eclipse.leshan.server.californium.observation.ObserveUtil; import org.eclipse.leshan.server.californium.registration.CaliforniumRegistrationStore; -import org.eclipse.leshan.server.redis.JedisLock; import org.eclipse.leshan.server.redis.RedisRegistrationStore; -import org.eclipse.leshan.server.redis.SingleInstanceJedisLock; import org.eclipse.leshan.server.redis.serialization.ObservationSerDes; import org.eclipse.leshan.server.redis.serialization.RegistrationSerDes; import org.eclipse.leshan.server.registration.Deregistration; @@ -38,11 +36,12 @@ import org.eclipse.leshan.server.registration.RegistrationUpdate; import org.eclipse.leshan.server.registration.UpdatedRegistration; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.springframework.data.redis.connection.RedisClusterConnection; +import org.springframework.data.redis.connection.RedisConnection; import org.springframework.data.redis.connection.RedisConnectionFactory; -import redis.clients.jedis.Jedis; -import redis.clients.jedis.ScanParams; -import redis.clients.jedis.ScanResult; -import redis.clients.jedis.Transaction; +import org.springframework.data.redis.core.Cursor; +import org.springframework.data.redis.core.ScanOptions; +import org.springframework.integration.redis.util.RedisLockRegistry; import java.net.InetSocketAddress; import java.util.ArrayList; @@ -50,13 +49,14 @@ import java.util.Arrays; import java.util.Collection; import java.util.Collections; import java.util.Iterator; +import java.util.LinkedList; import java.util.List; -import java.util.NoSuchElementException; import java.util.Set; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; +import java.util.concurrent.locks.Lock; import static java.nio.charset.StandardCharsets.UTF_8; @@ -92,7 +92,7 @@ public class TbLwM2mRedisRegistrationStore implements CaliforniumRegistrationSto private final int cleanLimit; // maximum number to clean in a clean period private final long gracePeriod; // in seconds - private final JedisLock lock; + private final RedisLockRegistry redisLock; public TbLwM2mRedisRegistrationStore(RedisConnectionFactory connectionFactory) { this(connectionFactory, DEFAULT_CLEAN_PERIOD, DEFAULT_GRACE_PERIOD, DEFAULT_CLEAN_LIMIT); // default clean period 60s @@ -106,20 +106,12 @@ public class TbLwM2mRedisRegistrationStore implements CaliforniumRegistrationSto public TbLwM2mRedisRegistrationStore(RedisConnectionFactory connectionFactory, ScheduledExecutorService schedExecutor, long cleanPeriodInSec, long lifetimeGracePeriodInSec, int cleanLimit) { - this(connectionFactory, schedExecutor, cleanPeriodInSec, lifetimeGracePeriodInSec, cleanLimit, new SingleInstanceJedisLock()); - } - - /** - * @since 1.1 - */ - public TbLwM2mRedisRegistrationStore(RedisConnectionFactory connectionFactory, ScheduledExecutorService schedExecutor, long cleanPeriodInSec, - long lifetimeGracePeriodInSec, int cleanLimit, JedisLock redisLock) { this.connectionFactory = connectionFactory; this.schedExecutor = schedExecutor; this.cleanPeriod = cleanPeriodInSec; this.cleanLimit = cleanLimit; this.gracePeriod = lifetimeGracePeriodInSec; - this.lock = redisLock; + this.redisLock = new RedisLockRegistry(connectionFactory, "Registration"); } /* *************** Redis Key utility function **************** */ @@ -135,76 +127,79 @@ public class TbLwM2mRedisRegistrationStore implements CaliforniumRegistrationSto return (prefix + registrationID).getBytes(); } - private byte[] toLockKey(String endpoint) { - return toKey(LOCK_EP, endpoint); + private String toLockKey(String endpoint) { + return new String(toKey(LOCK_EP, endpoint)); } - private byte[] toLockKey(byte[] endpoint) { - return toKey(LOCK_EP.getBytes(UTF_8), endpoint); + private String toLockKey(byte[] endpoint) { + return new String(toKey(LOCK_EP.getBytes(UTF_8), endpoint)); } /* *************** Leshan Registration API **************** */ @Override public Deregistration addRegistration(Registration registration) { - try (Jedis j = (Jedis) connectionFactory.getConnection().getNativeConnection()) { - byte[] lockValue = null; - byte[] lockKey = toLockKey(registration.getEndpoint()); + Lock lock = null; + try (var connection = connectionFactory.getConnection()) { + String lockKey = toLockKey(registration.getEndpoint()); try { - lockValue = lock.acquire(j, lockKey); - + lock = redisLock.obtain(lockKey); + lock.lock(); // add registration byte[] k = toEndpointKey(registration.getEndpoint()); - byte[] old = j.getSet(k, serializeReg(registration)); + byte[] old = connection.getSet(k, serializeReg(registration)); // add registration: secondary indexes byte[] regid_idx = toRegIdKey(registration.getId()); - j.set(regid_idx, registration.getEndpoint().getBytes(UTF_8)); + connection.set(regid_idx, registration.getEndpoint().getBytes(UTF_8)); byte[] addr_idx = toRegAddrKey(registration.getSocketAddress()); - j.set(addr_idx, registration.getEndpoint().getBytes(UTF_8)); + connection.set(addr_idx, registration.getEndpoint().getBytes(UTF_8)); // Add or update expiration - addOrUpdateExpiration(j, registration); + addOrUpdateExpiration(connection, registration); if (old != null) { Registration oldRegistration = deserializeReg(old); // remove old secondary index if (!registration.getId().equals(oldRegistration.getId())) - j.del(toRegIdKey(oldRegistration.getId())); + connection.del(toRegIdKey(oldRegistration.getId())); if (!oldRegistration.getSocketAddress().equals(registration.getSocketAddress())) { - removeAddrIndex(j, oldRegistration); + removeAddrIndex(connection, oldRegistration); } // remove old observation - Collection obsRemoved = unsafeRemoveAllObservations(j, oldRegistration.getId()); + Collection obsRemoved = unsafeRemoveAllObservations(connection, oldRegistration.getId()); return new Deregistration(oldRegistration, obsRemoved); } return null; } finally { - lock.release(j, lockKey, lockValue); + if (lock != null) { + lock.unlock(); + } } } } @Override public UpdatedRegistration updateRegistration(RegistrationUpdate update) { - try (Jedis j = (Jedis) connectionFactory.getConnection().getNativeConnection()) { + Lock lock = null; + try (var connection = connectionFactory.getConnection()) { // Fetch the registration ep by registration ID index - byte[] ep = j.get(toRegIdKey(update.getRegistrationId())); + byte[] ep = connection.get(toRegIdKey(update.getRegistrationId())); if (ep == null) { return null; } - byte[] lockValue = null; - byte[] lockKey = toLockKey(ep); + String lockKey = toLockKey(ep); try { - lockValue = lock.acquire(j, lockKey); + lock = redisLock.obtain(lockKey); + lock.lock(); // Fetch the registration - byte[] data = j.get(toEndpointKey(ep)); + byte[] data = connection.get(toEndpointKey(ep)); if (data == null) { return null; } @@ -214,40 +209,42 @@ public class TbLwM2mRedisRegistrationStore implements CaliforniumRegistrationSto Registration updatedRegistration = update.update(r); // Store the new registration - j.set(toEndpointKey(updatedRegistration.getEndpoint()), serializeReg(updatedRegistration)); + connection.set(toEndpointKey(updatedRegistration.getEndpoint()), serializeReg(updatedRegistration)); // Add or update expiration - addOrUpdateExpiration(j, updatedRegistration); + addOrUpdateExpiration(connection, updatedRegistration); // Update secondary index : // If registration is already associated to this address we don't care as we only want to keep the most // recent binding. byte[] addr_idx = toRegAddrKey(updatedRegistration.getSocketAddress()); - j.set(addr_idx, updatedRegistration.getEndpoint().getBytes(UTF_8)); + connection.set(addr_idx, updatedRegistration.getEndpoint().getBytes(UTF_8)); if (!r.getSocketAddress().equals(updatedRegistration.getSocketAddress())) { - removeAddrIndex(j, r); + removeAddrIndex(connection, r); } return new UpdatedRegistration(r, updatedRegistration); } finally { - lock.release(j, lockKey, lockValue); + if (lock != null) { + lock.unlock(); + } } } } @Override public Registration getRegistration(String registrationId) { - try (Jedis j = (Jedis) connectionFactory.getConnection().getNativeConnection()) { - return getRegistration(j, registrationId); + try (var connection = connectionFactory.getConnection()) { + return getRegistration(connection, registrationId); } } @Override public Registration getRegistrationByEndpoint(String endpoint) { Validate.notNull(endpoint); - try (Jedis j = (Jedis) connectionFactory.getConnection().getNativeConnection()) { - byte[] data = j.get(toEndpointKey(endpoint)); + try (var connection = connectionFactory.getConnection()) { + byte[] data = connection.get(toEndpointKey(endpoint)); if (data == null) { return null; } @@ -258,12 +255,12 @@ public class TbLwM2mRedisRegistrationStore implements CaliforniumRegistrationSto @Override public Registration getRegistrationByAdress(InetSocketAddress address) { Validate.notNull(address); - try (Jedis j = (Jedis) connectionFactory.getConnection().getNativeConnection()) { - byte[] ep = j.get(toRegAddrKey(address)); + try (var connection = connectionFactory.getConnection()) { + byte[] ep = connection.get(toRegAddrKey(address)); if (ep == null) { return null; } - byte[] data = j.get(toEndpointKey(ep)); + byte[] data = connection.get(toEndpointKey(ep)); if (data == null) { return null; } @@ -273,140 +270,99 @@ public class TbLwM2mRedisRegistrationStore implements CaliforniumRegistrationSto @Override public Iterator getAllRegistrations() { - return new TbLwM2mRedisRegistrationStore.RedisIterator(connectionFactory, new ScanParams().match(REG_EP + "*").count(100)); - } - - protected class RedisIterator implements Iterator { - - private final RedisConnectionFactory connectionFactory; - private final ScanParams scanParams; - - private String cursor; - private List scanResult; - - public RedisIterator(RedisConnectionFactory connectionFactory, ScanParams scanParams) { - this.connectionFactory = connectionFactory; - this.scanParams = scanParams; - // init scan result - scanNext("0"); - } - - private void scanNext(String cursor) { - try (Jedis j = (Jedis) connectionFactory.getConnection().getNativeConnection()) { - do { - ScanResult sr = j.scan(cursor.getBytes(), scanParams); - - this.scanResult = new ArrayList<>(); - if (sr.getResult() != null && !sr.getResult().isEmpty()) { - for (byte[] value : j.mget(sr.getResult().toArray(new byte[][]{}))) { - this.scanResult.add(deserializeReg(value)); - } - } - - cursor = sr.getCursor(); - } while (!"0".equals(cursor) && scanResult.isEmpty()); - - this.cursor = cursor; - } - } - - @Override - public boolean hasNext() { - if (!scanResult.isEmpty()) { - return true; - } - if ("0".equals(cursor)) { - // no more elements to scan - return false; - } - - // read more elements - scanNext(cursor); - return !scanResult.isEmpty(); - } - - @Override - public Registration next() { - if (!hasNext()) { - throw new NoSuchElementException(); + try (var connection = connectionFactory.getConnection()) { + Collection list = new LinkedList<>(); + ScanOptions scanOptions = ScanOptions.scanOptions().count(100).match(REG_EP + "*").build(); + List> scans = new ArrayList<>(); + if (connection instanceof RedisClusterConnection) { + ((RedisClusterConnection) connection).clusterGetNodes().forEach(node -> { + scans.add(((RedisClusterConnection) connection).scan(node, scanOptions)); + }); + } else { + scans.add(connection.scan(scanOptions)); } - return scanResult.remove(0); - } - @Override - public void remove() { - throw new UnsupportedOperationException(); + scans.forEach(scan -> { + scan.forEachRemaining(key -> { + byte[] element = connection.get(key); + list.add(deserializeReg(element)); + }); + }); + return list.iterator(); } } @Override public Deregistration removeRegistration(String registrationId) { - try (Jedis j = (Jedis) connectionFactory.getConnection().getNativeConnection()) { - return removeRegistration(j, registrationId, false); + try (var connection = connectionFactory.getConnection()) { + return removeRegistration(connection, registrationId, false); } } - private Deregistration removeRegistration(Jedis j, String registrationId, boolean removeOnlyIfNotAlive) { + private Deregistration removeRegistration(RedisConnection connection, String registrationId, boolean removeOnlyIfNotAlive) { // fetch the client ep by registration ID index - byte[] ep = j.get(toRegIdKey(registrationId)); + byte[] ep = connection.get(toRegIdKey(registrationId)); if (ep == null) { return null; } - byte[] lockValue = null; - byte[] lockKey = toLockKey(ep); + Lock lock = null; + String lockKey = toLockKey(ep); try { - lockValue = lock.acquire(j, lockKey); + lock = redisLock.obtain(lockKey); + lock.lock(); // fetch the client - byte[] data = j.get(toEndpointKey(ep)); + byte[] data = connection.get(toEndpointKey(ep)); if (data == null) { return null; } Registration r = deserializeReg(data); if (!removeOnlyIfNotAlive || !r.isAlive(gracePeriod)) { - long nbRemoved = j.del(toRegIdKey(r.getId())); + long nbRemoved = connection.del(toRegIdKey(r.getId())); if (nbRemoved > 0) { - j.del(toEndpointKey(r.getEndpoint())); - Collection obsRemoved = unsafeRemoveAllObservations(j, r.getId()); - removeAddrIndex(j, r); - removeExpiration(j, r); + connection.del(toEndpointKey(r.getEndpoint())); + Collection obsRemoved = unsafeRemoveAllObservations(connection, r.getId()); + removeAddrIndex(connection, r); + removeExpiration(connection, r); return new Deregistration(r, obsRemoved); } } return null; } finally { - lock.release(j, lockKey, lockValue); + if (lock != null) { + lock.unlock(); + } } } - private void removeAddrIndex(Jedis j, Registration registration) { + private void removeAddrIndex(RedisConnection connection, Registration registration) { // Watch the key to remove. byte[] regAddrKey = toRegAddrKey(registration.getSocketAddress()); - j.watch(regAddrKey); + connection.watch(regAddrKey); - byte[] epFromAddr = j.get(regAddrKey); + byte[] epFromAddr = connection.get(regAddrKey); // Delete the key if needed. if (Arrays.equals(epFromAddr, registration.getEndpoint().getBytes(UTF_8))) { // Try to delete the key - Transaction transaction = j.multi(); - transaction.del(regAddrKey); - transaction.exec(); + connection.multi(); + connection.del(regAddrKey); + connection.exec(); // if transaction failed this is not an issue as the socket address is probably reused and we don't neeed to // delete it anymore. } else { // the key must not be deleted. - j.unwatch(); + connection.unwatch(); } } - private void addOrUpdateExpiration(Jedis j, Registration registration) { - j.zadd(EXP_EP, registration.getExpirationTimeStamp(gracePeriod), registration.getEndpoint().getBytes(UTF_8)); + private void addOrUpdateExpiration(RedisConnection connection, Registration registration) { + connection.zAdd(EXP_EP, registration.getExpirationTimeStamp(gracePeriod), registration.getEndpoint().getBytes(UTF_8)); } - private void removeExpiration(Jedis j, Registration registration) { - j.zrem(EXP_EP, registration.getEndpoint().getBytes(UTF_8)); + private void removeExpiration(RedisConnection connection, Registration registration) { + connection.zRem(EXP_EP, registration.getEndpoint().getBytes(UTF_8)); } private byte[] toRegIdKey(String registrationId) { @@ -441,33 +397,35 @@ public class TbLwM2mRedisRegistrationStore implements CaliforniumRegistrationSto */ @Override public Collection addObservation(String registrationId, Observation observation) { - List removed = new ArrayList<>(); - try (Jedis j = (Jedis) connectionFactory.getConnection().getNativeConnection()) { + try (var connection = connectionFactory.getConnection()) { // fetch the client ep by registration ID index - byte[] ep = j.get(toRegIdKey(registrationId)); + byte[] ep = connection.get(toRegIdKey(registrationId)); if (ep == null) { return null; } - byte[] lockValue = null; - byte[] lockKey = toLockKey(ep); + Lock lock = null; + String lockKey = toLockKey(ep); try { - lockValue = lock.acquire(j, lockKey); + lock = redisLock.obtain(lockKey); + lock.lock(); // cancel existing observations for the same path and registration id. - for (Observation obs : getObservations(j, registrationId)) { + for (Observation obs : getObservations(connection, registrationId)) { if (observation.getPath().equals(obs.getPath()) && !Arrays.equals(observation.getId(), obs.getId())) { removed.add(obs); - unsafeRemoveObservation(j, registrationId, obs.getId()); + unsafeRemoveObservation(connection, registrationId, obs.getId()); } } } finally { - lock.release(j, lockKey, lockValue); + if (lock != null) { + lock.unlock(); + } } } return removed; @@ -475,29 +433,32 @@ public class TbLwM2mRedisRegistrationStore implements CaliforniumRegistrationSto @Override public Observation removeObservation(String registrationId, byte[] observationId) { - try (Jedis j = (Jedis) connectionFactory.getConnection().getNativeConnection()) { + try (var connection = connectionFactory.getConnection()) { // fetch the client ep by registration ID index - byte[] ep = j.get(toRegIdKey(registrationId)); + byte[] ep = connection.get(toRegIdKey(registrationId)); if (ep == null) { return null; } // remove observation - byte[] lockValue = null; - byte[] lockKey = toLockKey(ep); + Lock lock = null; + String lockKey = toLockKey(ep); try { - lockValue = lock.acquire(j, lockKey); + lock = redisLock.obtain(lockKey); + lock.lock(); Observation observation = build(get(new Token(observationId))); if (observation != null && registrationId.equals(observation.getRegistrationId())) { - unsafeRemoveObservation(j, registrationId, observationId); + unsafeRemoveObservation(connection, registrationId, observationId); return observation; } return null; } finally { - lock.release(j, lockKey, lockValue); + if (lock != null) { + lock.unlock(); + } } } } @@ -509,15 +470,15 @@ public class TbLwM2mRedisRegistrationStore implements CaliforniumRegistrationSto @Override public Collection getObservations(String registrationId) { - try (Jedis j = (Jedis) connectionFactory.getConnection().getNativeConnection()) { - return getObservations(j, registrationId); + try (var connection = connectionFactory.getConnection()) { + return getObservations(connection, registrationId); } } - private Collection getObservations(Jedis j, String registrationId) { + private Collection getObservations(RedisConnection connection, String registrationId) { Collection result = new ArrayList<>(); - for (byte[] token : j.lrange(toKey(OBS_TKNS_REGID_IDX, registrationId), 0, -1)) { - byte[] obs = j.get(toKey(OBS_TKN, token)); + for (byte[] token : connection.lRange(toKey(OBS_TKNS_REGID_IDX, registrationId), 0, -1)) { + byte[] obs = connection.get(toKey(OBS_TKN, token)); if (obs != null) { result.add(build(deserializeObs(obs))); } @@ -527,22 +488,24 @@ public class TbLwM2mRedisRegistrationStore implements CaliforniumRegistrationSto @Override public Collection removeObservations(String registrationId) { - try (Jedis j = (Jedis) connectionFactory.getConnection().getNativeConnection()) { + try (var connection = connectionFactory.getConnection()) { // check registration exists - Registration registration = getRegistration(j, registrationId); + Registration registration = getRegistration(connection, registrationId); if (registration == null) return Collections.emptyList(); // get endpoint and create lock String endpoint = registration.getEndpoint(); - byte[] lockValue = null; - byte[] lockKey = toKey(LOCK_EP, endpoint); + Lock lock = null; + String lockKey = toLockKey(endpoint); try { - lockValue = lock.acquire(j, lockKey); - - return unsafeRemoveAllObservations(j, registrationId); + lock = redisLock.obtain(lockKey); + lock.lock(); + return unsafeRemoveAllObservations(connection, registrationId); } finally { - lock.release(j, lockKey, lockValue); + if (lock != null) { + lock.unlock(); + } } } } @@ -565,31 +528,32 @@ public class TbLwM2mRedisRegistrationStore implements CaliforniumRegistrationSto String endpoint = ObserveUtil.validateCoapObservation(obs); org.eclipse.californium.core.observe.Observation previousObservation = null; - try (Jedis j = (Jedis) connectionFactory.getConnection().getNativeConnection()) { - byte[] lockValue = null; - byte[] lockKey = toKey(LOCK_EP, endpoint); + try (var connection = connectionFactory.getConnection()) { + Lock lock = null; + String lockKey = toLockKey(endpoint); try { - lockValue = lock.acquire(j, lockKey); + lock = redisLock.obtain(lockKey); + lock.lock(); String registrationId = ObserveUtil.extractRegistrationId(obs); - if (!j.exists(toRegIdKey(registrationId))) + if (!connection.exists(toRegIdKey(registrationId))) throw new ObservationStoreException("no registration for this Id"); byte[] key = toKey(OBS_TKN, obs.getRequest().getToken().getBytes()); byte[] serializeObs = serializeObs(obs); byte[] previousValue; if (ifAbsent) { - previousValue = j.get(key); + previousValue = connection.get(key); if (previousValue == null || previousValue.length == 0) { - j.set(key, serializeObs); + connection.set(key, serializeObs); } else { return deserializeObs(previousValue); } } else { - previousValue = j.getSet(key, serializeObs); + previousValue = connection.getSet(key, serializeObs); } // secondary index to get the list by registrationId - j.lpush(toKey(OBS_TKNS_REGID_IDX, registrationId), obs.getRequest().getToken().getBytes()); + connection.lPush(toKey(OBS_TKNS_REGID_IDX, registrationId), obs.getRequest().getToken().getBytes()); // log any collisions if (previousValue != null && previousValue.length != 0) { @@ -599,7 +563,9 @@ public class TbLwM2mRedisRegistrationStore implements CaliforniumRegistrationSto previousObservation.getRequest(), obs.getRequest()); } } finally { - lock.release(j, lockKey, lockValue); + if (lock != null) { + lock.unlock(); + } } } return previousObservation; @@ -607,17 +573,17 @@ public class TbLwM2mRedisRegistrationStore implements CaliforniumRegistrationSto @Override public void remove(Token token) { - try (Jedis j = (Jedis) connectionFactory.getConnection().getNativeConnection()) { + try (var connection = connectionFactory.getConnection()) { byte[] tokenKey = toKey(OBS_TKN, token.getBytes()); // fetch the observation by token - byte[] serializedObs = j.get(tokenKey); + byte[] serializedObs = connection.get(tokenKey); if (serializedObs == null) return; org.eclipse.californium.core.observe.Observation obs = deserializeObs(serializedObs); String registrationId = ObserveUtil.extractRegistrationId(obs); - Registration registration = getRegistration(j, registrationId); + Registration registration = getRegistration(connection, registrationId); if (registration == null) { LOG.warn("Unable to remove observation {}, registration {} does not exist anymore", obs.getRequest(), registrationId); @@ -625,14 +591,17 @@ public class TbLwM2mRedisRegistrationStore implements CaliforniumRegistrationSto } String endpoint = registration.getEndpoint(); - byte[] lockValue = null; - byte[] lockKey = toKey(LOCK_EP, endpoint); + Lock lock = null; + String lockKey = toLockKey(endpoint); try { - lockValue = lock.acquire(j, lockKey); + lock = redisLock.obtain(lockKey); + lock.lock(); - unsafeRemoveObservation(j, registrationId, token.getBytes()); + unsafeRemoveObservation(connection, registrationId, token.getBytes()); } finally { - lock.release(j, lockKey, lockValue); + if (lock != null) { + lock.unlock(); + } } } @@ -640,8 +609,8 @@ public class TbLwM2mRedisRegistrationStore implements CaliforniumRegistrationSto @Override public org.eclipse.californium.core.observe.Observation get(Token token) { - try (Jedis j = (Jedis) connectionFactory.getConnection().getNativeConnection()) { - byte[] obs = j.get(toKey(OBS_TKN, token.getBytes())); + try (var connection = connectionFactory.getConnection()) { + byte[] obs = connection.get(toKey(OBS_TKN, token.getBytes())); if (obs == null) { return null; } else { @@ -652,12 +621,12 @@ public class TbLwM2mRedisRegistrationStore implements CaliforniumRegistrationSto /* *************** Observation utility functions **************** */ - private Registration getRegistration(Jedis j, String registrationId) { - byte[] ep = j.get(toRegIdKey(registrationId)); + private Registration getRegistration(RedisConnection connection, String registrationId) { + byte[] ep = connection.get(toRegIdKey(registrationId)); if (ep == null) { return null; } - byte[] data = j.get(toEndpointKey(ep)); + byte[] data = connection.get(toEndpointKey(ep)); if (data == null) { return null; } @@ -665,25 +634,25 @@ public class TbLwM2mRedisRegistrationStore implements CaliforniumRegistrationSto return deserializeReg(data); } - private void unsafeRemoveObservation(Jedis j, String registrationId, byte[] observationId) { - if (j.del(toKey(OBS_TKN, observationId)) > 0L) { - j.lrem(toKey(OBS_TKNS_REGID_IDX, registrationId), 0, observationId); + private void unsafeRemoveObservation(RedisConnection connection, String registrationId, byte[] observationId) { + if (connection.del(toKey(OBS_TKN, observationId)) > 0L) { + connection.lRem(toKey(OBS_TKNS_REGID_IDX, registrationId), 0, observationId); } } - private Collection unsafeRemoveAllObservations(Jedis j, String registrationId) { + private Collection unsafeRemoveAllObservations(RedisConnection connection, String registrationId) { Collection removed = new ArrayList<>(); byte[] regIdKey = toKey(OBS_TKNS_REGID_IDX, registrationId); // fetch all observations by token - for (byte[] token : j.lrange(regIdKey, 0, -1)) { - byte[] obs = j.get(toKey(OBS_TKN, token)); + for (byte[] token : connection.lRange(regIdKey, 0, -1)) { + byte[] obs = connection.get(toKey(OBS_TKN, token)); if (obs != null) { removed.add(build(deserializeObs(obs))); } - j.del(toKey(OBS_TKN, token)); + connection.del(toKey(OBS_TKN, token)); } - j.del(regIdKey); + connection.del(regIdKey); return removed; } @@ -754,14 +723,14 @@ public class TbLwM2mRedisRegistrationStore implements CaliforniumRegistrationSto @Override public void run() { - try (Jedis j = (Jedis) connectionFactory.getConnection().getNativeConnection()) { - Set endpointsExpired = j.zrangeByScore(EXP_EP, Double.NEGATIVE_INFINITY, + try (var connection = connectionFactory.getConnection()) { + Set endpointsExpired = connection.zRangeByScore(EXP_EP, Double.NEGATIVE_INFINITY, System.currentTimeMillis(), 0, cleanLimit); for (byte[] endpoint : endpointsExpired) { - Registration r = deserializeReg(j.get(toEndpointKey(endpoint))); + Registration r = deserializeReg(connection.get(toEndpointKey(endpoint))); if (!r.isAlive(gracePeriod)) { - Deregistration dereg = removeRegistration(j, r.getId(), true); + Deregistration dereg = removeRegistration(connection, r.getId(), true); if (dereg != null) expirationListener.registrationExpired(dereg.getRegistration(), dereg.getObservations()); } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mRedisSecurityStore.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mRedisSecurityStore.java index 47b13bb72b..4cfe2a6829 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mRedisSecurityStore.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mRedisSecurityStore.java @@ -20,13 +20,15 @@ import org.eclipse.leshan.server.security.EditableSecurityStore; import org.eclipse.leshan.server.security.NonUniqueSecurityInfoException; import org.eclipse.leshan.server.security.SecurityInfo; import org.eclipse.leshan.server.security.SecurityStoreListener; +import org.springframework.data.redis.connection.RedisClusterConnection; import org.springframework.data.redis.connection.RedisConnectionFactory; -import redis.clients.jedis.Jedis; -import redis.clients.jedis.ScanParams; -import redis.clients.jedis.ScanResult; +import org.springframework.data.redis.core.Cursor; +import org.springframework.data.redis.core.ScanOptions; +import java.util.ArrayList; import java.util.Collection; import java.util.LinkedList; +import java.util.List; public class TbLwM2mRedisSecurityStore implements EditableSecurityStore { private static final String SEC_EP = "SEC#EP#"; @@ -42,8 +44,8 @@ public class TbLwM2mRedisSecurityStore implements EditableSecurityStore { @Override public SecurityInfo getByEndpoint(String endpoint) { - try (Jedis j = (Jedis) connectionFactory.getConnection().getNativeConnection()) { - byte[] data = j.get((SEC_EP + endpoint).getBytes()); + try (var connection = connectionFactory.getConnection()) { + byte[] data = connection.get((SEC_EP + endpoint).getBytes()); if (data == null) { return null; } else { @@ -54,12 +56,12 @@ public class TbLwM2mRedisSecurityStore implements EditableSecurityStore { @Override public SecurityInfo getByIdentity(String identity) { - try (Jedis j = (Jedis) connectionFactory.getConnection().getNativeConnection()) { - String ep = j.hget(PSKID_SEC, identity); + try (var connection = connectionFactory.getConnection()) { + byte[] ep = connection.hGet(PSKID_SEC.getBytes(), identity.getBytes()); if (ep == null) { return null; } else { - byte[] data = j.get((SEC_EP + ep).getBytes()); + byte[] data = connection.get((SEC_EP + new String(ep)).getBytes()); if (data == null) { return null; } else { @@ -71,18 +73,24 @@ public class TbLwM2mRedisSecurityStore implements EditableSecurityStore { @Override public Collection getAll() { - try (Jedis j = (Jedis) connectionFactory.getConnection().getNativeConnection()) { - ScanParams params = new ScanParams().match(SEC_EP + "*").count(100); + try (var connection = connectionFactory.getConnection()) { Collection list = new LinkedList<>(); - String cursor = "0"; - do { - ScanResult res = j.scan(cursor.getBytes(), params); - for (byte[] key : res.getResult()) { - byte[] element = j.get(key); + ScanOptions scanOptions = ScanOptions.scanOptions().count(100).match(SEC_EP + "*").build(); + List> scans = new ArrayList<>(); + if (connection instanceof RedisClusterConnection) { + ((RedisClusterConnection) connection).clusterGetNodes().forEach(node -> { + scans.add(((RedisClusterConnection) connection).scan(node, scanOptions)); + }); + } else { + scans.add(connection.scan(scanOptions)); + } + + scans.forEach(scan -> { + scan.forEachRemaining(key -> { + byte[] element = connection.get(key); list.add(deserialize(element)); - } - cursor = res.getCursor(); - } while (!"0".equals(cursor)); + }); + }); return list; } } @@ -90,21 +98,21 @@ public class TbLwM2mRedisSecurityStore implements EditableSecurityStore { @Override public SecurityInfo add(SecurityInfo info) throws NonUniqueSecurityInfoException { byte[] data = serialize(info); - try (Jedis j = (Jedis) connectionFactory.getConnection().getNativeConnection()) { + try (var connection = connectionFactory.getConnection()) { if (info.getIdentity() != null) { // populate the secondary index (security info by PSK id) - String oldEndpoint = j.hget(PSKID_SEC, info.getIdentity()); - if (oldEndpoint != null && !oldEndpoint.equals(info.getEndpoint())) { + String oldEndpoint = new String(connection.hGet(PSKID_SEC.getBytes(), info.getIdentity().getBytes())); + if (!oldEndpoint.equals(info.getEndpoint())) { throw new NonUniqueSecurityInfoException("PSK Identity " + info.getIdentity() + " is already used"); } - j.hset(PSKID_SEC.getBytes(), info.getIdentity().getBytes(), info.getEndpoint().getBytes()); + connection.hSet(PSKID_SEC.getBytes(), info.getIdentity().getBytes(), info.getEndpoint().getBytes()); } - byte[] previousData = j.getSet((SEC_EP + info.getEndpoint()).getBytes(), data); + byte[] previousData = connection.getSet((SEC_EP + info.getEndpoint()).getBytes(), data); SecurityInfo previous = previousData == null ? null : deserialize(previousData); String previousIdentity = previous == null ? null : previous.getIdentity(); if (previousIdentity != null && !previousIdentity.equals(info.getIdentity())) { - j.hdel(PSKID_SEC, previousIdentity); + connection.hDel(PSKID_SEC.getBytes(), previousIdentity.getBytes()); } return previous; @@ -113,15 +121,15 @@ public class TbLwM2mRedisSecurityStore implements EditableSecurityStore { @Override public SecurityInfo remove(String endpoint, boolean infosAreCompromised) { - try (Jedis j = (Jedis) connectionFactory.getConnection().getNativeConnection()) { - byte[] data = j.get((SEC_EP + endpoint).getBytes()); + try (var connection = connectionFactory.getConnection()) { + byte[] data = connection.get((SEC_EP + endpoint).getBytes()); if (data != null) { SecurityInfo info = deserialize(data); if (info.getIdentity() != null) { - j.hdel(PSKID_SEC.getBytes(), info.getIdentity().getBytes()); + connection.hDel(PSKID_SEC.getBytes(), info.getIdentity().getBytes()); } - j.del((SEC_EP + endpoint).getBytes()); + connection.del((SEC_EP + endpoint).getBytes()); if (listener != null) { listener.securityInfoRemoved(infosAreCompromised, info); } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mStoreFactory.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mStoreFactory.java index 1fcbd15e33..164a78b0a3 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mStoreFactory.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mStoreFactory.java @@ -26,6 +26,7 @@ import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Component; import org.thingsboard.server.cache.TBRedisCacheConfiguration; import org.thingsboard.server.queue.util.TbLwM2mTransportComponent; +import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig; import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClientContext; import java.util.Optional; @@ -37,6 +38,9 @@ public class TbLwM2mStoreFactory { @Autowired(required = false) private Optional redisConfiguration; + @Autowired + private LwM2MTransportServerConfig config; + @Autowired @Lazy private LwM2mClientContext clientContext; @@ -47,7 +51,7 @@ public class TbLwM2mStoreFactory { @Bean private CaliforniumRegistrationStore registrationStore() { return redisConfiguration.isPresent() && useRedis ? - new TbLwM2mRedisRegistrationStore(redisConfiguration.get().redisConnectionFactory()) : new InMemoryRegistrationStore(); + new TbLwM2mRedisRegistrationStore(redisConfiguration.get().redisConnectionFactory()) : new InMemoryRegistrationStore(config.getCleanPeriodInSec()); } @Bean diff --git a/pom.xml b/pom.xml index 4073fe824e..29a373f2e4 100755 --- a/pom.xml +++ b/pom.xml @@ -1430,6 +1430,11 @@ spring-data-redis ${spring-data-redis.version} + + org.springframework.integration + spring-integration-redis + ${spring.version} + redis.clients jedis diff --git a/transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml b/transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml index 1de3ffb630..39acda4590 100644 --- a/transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml +++ b/transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml @@ -98,7 +98,6 @@ transport: # Enable/disable http/mqtt/coap transport protocols (has higher priority than certain protocol's 'enabled' property) api_enabled: "${TB_TRANSPORT_API_ENABLED:true}" # Local LwM2M transport parameters - # Local LwM2M transport parameters lwm2m: # Enable/disable lvm2m transport protocol. enabled: "${LWM2M_ENABLED:true}" @@ -144,6 +143,8 @@ transport: recommended_supported_groups: "${LWM2M_RECOMMENDED_SUPPORTED_GROUPS:true}" response_pool_size: "${LWM2M_RESPONSE_POOL_SIZE:100}" registered_pool_size: "${LWM2M_REGISTERED_POOL_SIZE:10}" + registration_store_pool_size: "${LWM2M_REGISTRATION_STORE_POOL_SIZE:100}" + clean_period_in_sec: "${LWM2M_CLEAN_PERIOD_IN_SEC:2}" update_registered_pool_size: "${LWM2M_UPDATE_REGISTERED_POOL_SIZE:10}" un_registered_pool_size: "${LWM2M_UN_REGISTERED_POOL_SIZE:10}" log_max_length: "${LWM2M_LOG_MAX_LENGTH:100}"