From 992767d7b102e3560e3682e4630a29c5843ede01 Mon Sep 17 00:00:00 2001 From: nick Date: Thu, 28 Dec 2023 16:44:31 +0200 Subject: [PATCH] lwm2m: redis --- .../store/TbLwM2mRedisRegistrationStore.java | 92 +++++++++---------- 1 file changed, 41 insertions(+), 51 deletions(-) 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 40a281728b..fa44c2a93b 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 @@ -333,12 +333,50 @@ public class TbLwM2mRedisRegistrationStore implements RegistrationStore, Startab @Override public Collection addObservation(String registrationId, Observation observation, boolean addIfAbsent) { - return null; + List removed = new ArrayList<>(); + try (var connection = connectionFactory.getConnection()) { + + // fetch the client ep by registration ID index + byte[] ep = connection.get(toRegIdKey(registrationId)); + if (ep == null) { + return null; + } + + Lock lock = null; + String lockKey = toLockKey(ep); + + try { + lock = redisLock.obtain(lockKey); + lock.lock(); + + // cancel existing observations for the same path and registration id. + for (Observation obs : getObservations(connection, registrationId)) { + //TODO: should be able to use CompositeObservation + if (((SingleObservation)observation).getPath().equals(((SingleObservation)obs).getPath()) + && !observation.getId().equals(obs.getId())) { + removed.add(obs); + unsafeRemoveObservation(connection, registrationId, obs.getId().getBytes()); + } + } + + } finally { + if (lock != null) { + lock.unlock(); + } + } + } + return removed; } + @Override + public Collection getObservations(String registrationId) { + try (var connection = connectionFactory.getConnection()) { + return getObservations(connection, registrationId); + } + } @Override public Observation getObservation(String registrationId, ObservationIdentifier observationId) { - return null; + return getObservations(registrationId).stream().filter(o -> o.getId()==observationId).findFirst().get(); } @Override @@ -348,7 +386,7 @@ public class TbLwM2mRedisRegistrationStore implements RegistrationStore, Startab @Override public Observation removeObservation(String registrationId, ObservationIdentifier observationId) { - return null; + return removeObservation(registrationId, observationId.getBytes()); } private Deregistration removeRegistration(RedisConnection connection, String registrationId, boolean removeOnlyIfNotAlive) { @@ -461,42 +499,6 @@ public class TbLwM2mRedisRegistrationStore implements RegistrationStore, Startab * org.eclipse.californium.core.observe.ObservationStore#add method) */ - public Collection addObservation(String registrationId, Observation observation) { - List removed = new ArrayList<>(); - try (var connection = connectionFactory.getConnection()) { - - // fetch the client ep by registration ID index - byte[] ep = connection.get(toRegIdKey(registrationId)); - if (ep == null) { - return null; - } - - Lock lock = null; - String lockKey = toLockKey(ep); - - try { - lock = redisLock.obtain(lockKey); - lock.lock(); - - // cancel existing observations for the same path and registration id. - for (Observation obs : getObservations(connection, registrationId)) { - //TODO: should be able to use CompositeObservation - if (((SingleObservation)observation).getPath().equals(((SingleObservation)obs).getPath()) - && !observation.getId().equals(obs.getId())) { - removed.add(obs); - unsafeRemoveObservation(connection, registrationId, obs.getId().getBytes()); - } - } - - } finally { - if (lock != null) { - lock.unlock(); - } - } - } - return removed; - } - public Observation removeObservation(String registrationId, byte[] observationId) { try (var connection = connectionFactory.getConnection()) { @@ -529,18 +531,6 @@ public class TbLwM2mRedisRegistrationStore implements RegistrationStore, Startab } } - public Observation getObservation(String registrationId, byte[] observationId) { -// return build(get(new Token(observationId))); - return build(null); - } - - @Override - public Collection getObservations(String registrationId) { - try (var connection = connectionFactory.getConnection()) { - return getObservations(connection, registrationId); - } - } - private Collection getObservations(RedisConnection connection, String registrationId) { Collection result = new ArrayList<>(); for (byte[] token : connection.lRange(toKey(OBS_TKNS_REGID_IDX, registrationId), 0, -1)) {