mirror of
https://github.com/eclipse-sailing-analytics/sailing-analytics.git
synced 2026-10-10 22:30:56 +00:00
Merge branch 'bug3875'
This commit is contained in:
commit
79ce34721c
7 files changed
+263
-124
No files matched your search
+25
-14
@@ -37,6 +37,8 @@ import com.sap.sse.common.TypeBasedServiceFinder;
|
||||
import com.sap.sse.common.TypeBasedServiceFinderFactory;
|
||||
import com.sap.sse.common.Util;
|
||||
import com.sap.sse.common.impl.TimeRangeImpl;
|
||||
import com.sap.sse.concurrent.LockUtil;
|
||||
import com.sap.sse.concurrent.NamedReentrantReadWriteLock;
|
||||
|
||||
/**
|
||||
* At the moment, the timerange covered by the fixes for a device, and the number of fixes for a device are stored in a
|
||||
@@ -52,6 +54,10 @@ public class MongoSensorFixStoreImpl implements MongoSensorFixStore {
|
||||
private final DBCollection fixesCollection;
|
||||
private final DBCollection metadataCollection;
|
||||
private final MongoObjectFactoryImpl mongoOF;
|
||||
/**
|
||||
* Lock object to be used when accessing {@link #listeners}.
|
||||
*/
|
||||
private final NamedReentrantReadWriteLock listenersLock = new NamedReentrantReadWriteLock("Listeners collection lock of " + MongoSensorFixStoreImpl.class.getName(), false);
|
||||
private final Map<DeviceIdentifier, Set<FixReceivedListener<? extends Timed>>> listeners = new HashMap<>();
|
||||
|
||||
public MongoSensorFixStoreImpl(MongoObjectFactory mongoObjectFactory, DomainObjectFactory domainObjectFactory,
|
||||
@@ -168,9 +174,7 @@ public class MongoSensorFixStoreImpl implements MongoSensorFixStore {
|
||||
logger.log(Level.WARNING, "Could not store fix in MongoDB");
|
||||
e.printStackTrace();
|
||||
}
|
||||
for (FixT fix : fixes) {
|
||||
notifyListeners(device, fix);
|
||||
}
|
||||
notifyListeners(device, fixes);
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -179,23 +183,30 @@ public class MongoSensorFixStoreImpl implements MongoSensorFixStore {
|
||||
}
|
||||
|
||||
@SuppressWarnings({ "unchecked", "rawtypes" })
|
||||
private <FixT extends Timed> void notifyListeners(DeviceIdentifier device, FixT fix) {
|
||||
for (FixReceivedListener<FixT> listener : Util.<DeviceIdentifier, Set<FixReceivedListener<FixT>>> get(
|
||||
(Map) listeners, device, Collections.emptySet())) {
|
||||
listener.fixReceived(device, fix);
|
||||
}
|
||||
private <FixT extends Timed> void notifyListeners(DeviceIdentifier device, Iterable<FixT> fixes) {
|
||||
LockUtil.executeWithReadLock(listenersLock, () -> {
|
||||
for (FixT fix : fixes) {
|
||||
for (FixReceivedListener<FixT> listener : Util.<DeviceIdentifier, Set<FixReceivedListener<FixT>>> get(
|
||||
(Map) listeners, device, Collections.emptySet())) {
|
||||
listener.fixReceived(device, fix);
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
@Override
|
||||
public synchronized void addListener(FixReceivedListener<? extends Timed> listener, DeviceIdentifier device) {
|
||||
Util.addToValueSet(listeners, device, listener);
|
||||
public void addListener(FixReceivedListener<? extends Timed> listener, DeviceIdentifier device) {
|
||||
LockUtil.executeWithWriteLock(listenersLock, () -> Util.addToValueSet(listeners, device, listener));
|
||||
}
|
||||
|
||||
@Override
|
||||
public synchronized void removeListener(FixReceivedListener<? extends Timed> listener) {
|
||||
for (Set<FixReceivedListener<? extends Timed>> set : listeners.values()) {
|
||||
set.remove(listener);
|
||||
}
|
||||
public void removeListener(FixReceivedListener<? extends Timed> listener) {
|
||||
LockUtil.executeWithWriteLock(listenersLock, () -> Util.removeFromAllValueSets(listeners, listener));
|
||||
}
|
||||
|
||||
@Override
|
||||
public void removeListener(FixReceivedListener<? extends Timed> listener, DeviceIdentifier device) {
|
||||
LockUtil.executeWithWriteLock(listenersLock, () -> Util.removeFromValueSet(listeners, device, listener));
|
||||
}
|
||||
|
||||
private DBObject getDeviceQuery(DeviceIdentifier device)
|
||||
|
||||
+85
@@ -0,0 +1,85 @@
|
||||
package com.sap.sailing.domain.racelogtracking.test.impl;
|
||||
|
||||
import java.util.concurrent.CyclicBarrier;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.TimeoutException;
|
||||
|
||||
import org.junit.Rule;
|
||||
import org.junit.Test;
|
||||
import org.junit.rules.Timeout;
|
||||
|
||||
import com.sap.sailing.domain.common.tracking.GPSFixMoving;
|
||||
import com.sap.sailing.domain.persistence.racelog.tracking.impl.MongoSensorFixStoreImpl;
|
||||
import com.sap.sailing.domain.racelog.tracking.FixReceivedListener;
|
||||
import com.sap.sailing.domain.racelogtracking.DeviceIdentifier;
|
||||
import com.sap.sailing.domain.racelogtracking.test.AbstractGPSFixStoreTest;
|
||||
|
||||
public class GPSFixStoreListenerTest extends AbstractGPSFixStoreTest {
|
||||
@Rule
|
||||
public Timeout GPSFixStoreListenerTestTimeout = new Timeout(3 * 1000);
|
||||
|
||||
/**
|
||||
* {@link MongoSensorFixStoreImpl} had broken synchronization of the listeners collection (add/removeListener
|
||||
* methods were synchronized but notifyListeners was not synchronized).
|
||||
*/
|
||||
@Test(expected = TimoutRuntimeException.class)
|
||||
public void lockingOfGPSFixStoreListenersIsWorkingCorrectly() throws InterruptedException {
|
||||
CyclicBarrier barrier = new CyclicBarrier(2);
|
||||
// We need 3 listener instances to guarantee that the iterator isn't finished
|
||||
// when adding another listener in the thread below.
|
||||
store.addListener(new ListenerAwaitingBarier(barrier), device);
|
||||
store.addListener(new ListenerAwaitingBarier(barrier), device);
|
||||
store.addListener(new ListenerAwaitingBarier(barrier), device);
|
||||
|
||||
Thread thread = new Thread() {
|
||||
public void run() {
|
||||
try {
|
||||
barrier.await(100, TimeUnit.MILLISECONDS);
|
||||
// During iteration in the main thread this causes a modification that makes the iterator throw a
|
||||
// ConcurrentModificationException on next()
|
||||
store.addListener((DeviceIdentifier device, GPSFixMoving fix) -> {}, device);
|
||||
barrier.await(100, TimeUnit.MILLISECONDS);
|
||||
barrier.await(100, TimeUnit.MILLISECONDS);
|
||||
} catch (Exception e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
};
|
||||
};
|
||||
thread.start();
|
||||
try {
|
||||
store.storeFix(device, createFix(100, 10, 20, 30, 40));
|
||||
} finally {
|
||||
// This ensures that the thread is terminated when the test finishes
|
||||
// JUnit may behave crazy if there are additional tests running after the test finished
|
||||
thread.join(500);
|
||||
}
|
||||
}
|
||||
|
||||
private static class ListenerAwaitingBarier implements FixReceivedListener<GPSFixMoving> {
|
||||
|
||||
private final CyclicBarrier barrier;
|
||||
|
||||
public ListenerAwaitingBarier(CyclicBarrier barrier) {
|
||||
this.barrier = barrier;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void fixReceived(DeviceIdentifier device, GPSFixMoving fix) {
|
||||
try {
|
||||
barrier.await(100, TimeUnit.MILLISECONDS);
|
||||
} catch (TimeoutException e) {
|
||||
throw new TimoutRuntimeException(e);
|
||||
} catch (Exception e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private static class TimoutRuntimeException extends RuntimeException {
|
||||
private static final long serialVersionUID = 6349762933223278846L;
|
||||
|
||||
public TimoutRuntimeException(TimeoutException e) {
|
||||
super(e);
|
||||
}
|
||||
}
|
||||
}
|
||||
+22
-19
@@ -43,8 +43,6 @@ import com.sap.sse.common.Timed;
|
||||
import com.sap.sse.common.Util;
|
||||
import com.sap.sse.common.WithID;
|
||||
import com.sap.sse.common.impl.TimeRangeImpl;
|
||||
import com.sap.sse.concurrent.LockUtil;
|
||||
import com.sap.sse.concurrent.NamedReentrantReadWriteLock;
|
||||
|
||||
/**
|
||||
* This class listens to RaceLog Events, changes to the race and fix loading events and properly handles mappings and
|
||||
@@ -60,7 +58,6 @@ import com.sap.sse.concurrent.NamedReentrantReadWriteLock;
|
||||
public class FixLoaderAndTracker implements TrackingDataLoader {
|
||||
private static final Logger logger = Logger.getLogger(FixLoaderAndTracker.class.getName());
|
||||
protected final DynamicTrackedRace trackedRace;
|
||||
private final NamedReentrantReadWriteLock loadingFromFixStoreLock;
|
||||
private final SensorFixStore sensorFixStore;
|
||||
private final GPSFixStore gpsFixStore;
|
||||
private RegattaLogDeviceMappings<WithID> deviceMappings;
|
||||
@@ -158,8 +155,6 @@ public class FixLoaderAndTracker implements TrackingDataLoader {
|
||||
this.gpsFixStore = new GPSFixStoreImpl(sensorFixStore);
|
||||
this.sensorFixMapperFactory = sensorFixMapperFactory;
|
||||
this.trackedRace = trackedRace;
|
||||
loadingFromFixStoreLock = new NamedReentrantReadWriteLock(
|
||||
"Loading from SensorFix store lock for tracked race " + trackedRace.getRace().getName(), false);
|
||||
|
||||
startTracking();
|
||||
}
|
||||
@@ -254,7 +249,8 @@ public class FixLoaderAndTracker implements TrackingDataLoader {
|
||||
|
||||
protected void startTracking() {
|
||||
trackedRace.addListener(raceChangeListener);
|
||||
this.deviceMappings = new FixLoaderDeviceMappings(trackedRace.getAttachedRegattaLogs());
|
||||
this.deviceMappings = new FixLoaderDeviceMappings(trackedRace.getAttachedRegattaLogs(),
|
||||
trackedRace.getRace().getName());
|
||||
}
|
||||
|
||||
private synchronized void waitForLoadingToFinishRunning() {
|
||||
@@ -287,7 +283,16 @@ public class FixLoaderAndTracker implements TrackingDataLoader {
|
||||
});
|
||||
}
|
||||
|
||||
private void updateConcurrent(final Runnable updateCallback) {
|
||||
/**
|
||||
* This method runs the given update callback in a separate {@link Thread} by handling technical concurrency aspects
|
||||
* and potential {@link #preemptiveStopRequested preemptive stop requests} internally. Thus, it separates the
|
||||
* functional updating process from technical aspects.
|
||||
*
|
||||
* @param updateCallback
|
||||
* the {@link Runnable} callback used to run the update
|
||||
*/
|
||||
// TODO Consider using a thread pool here, after merging bug3864 into master
|
||||
private void updateAsyncInternal(final Runnable updateCallback) {
|
||||
synchronized (FixLoaderAndTracker.this) {
|
||||
activeLoaders.incrementAndGet();
|
||||
setStatusAndProgress(TrackedRaceStatusEnum.LOADING, 0.5);
|
||||
@@ -300,17 +305,14 @@ public class FixLoaderAndTracker implements TrackingDataLoader {
|
||||
if (!preemptiveStopRequested.get()) {
|
||||
trackedRace.lockForSerializationRead();
|
||||
setStatusAndProgress(TrackedRaceStatusEnum.LOADING, 0.5);
|
||||
LockUtil.lockForWrite(loadingFromFixStoreLock);
|
||||
synchronized (FixLoaderAndTracker.this) {
|
||||
FixLoaderAndTracker.this.notifyAll();
|
||||
}
|
||||
updateCallback.run();
|
||||
}
|
||||
} finally {
|
||||
LockUtil.unlockAfterWrite(loadingFromFixStoreLock);
|
||||
synchronized (FixLoaderAndTracker.this) {
|
||||
int currentActiveLoaders;
|
||||
currentActiveLoaders = activeLoaders.decrementAndGet();
|
||||
int currentActiveLoaders = activeLoaders.decrementAndGet();
|
||||
FixLoaderAndTracker.this.notifyAll();
|
||||
if (currentActiveLoaders == 0) {
|
||||
setStatusAndProgress(stopRequested.get() ? TrackedRaceStatusEnum.FINISHED
|
||||
@@ -338,27 +340,28 @@ public class FixLoaderAndTracker implements TrackingDataLoader {
|
||||
}
|
||||
|
||||
private class FixLoaderDeviceMappings extends RegattaLogDeviceMappings<WithID> {
|
||||
public FixLoaderDeviceMappings(Iterable<RegattaLog> initialRegattaLogs) {
|
||||
super(initialRegattaLogs);
|
||||
public FixLoaderDeviceMappings(Iterable<RegattaLog> initialRegattaLogs, String raceNameForLock) {
|
||||
super(initialRegattaLogs, raceNameForLock);
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void updateMappings() {
|
||||
updateConcurrent(() -> {
|
||||
FixLoaderDeviceMappings.super.updateMappings();
|
||||
// add listeners for devices in mappings already present
|
||||
forEachDevice((device) -> sensorFixStore.addListener(listener, device));
|
||||
});
|
||||
updateAsyncInternal(FixLoaderDeviceMappings.super::updateMappings);
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void mappingRemoved(DeviceMappingWithRegattaLogEvent<WithID> mapping) {
|
||||
final DeviceIdentifier device = mapping.getDevice();
|
||||
if(!hasMappingForDevice(device)) {
|
||||
sensorFixStore.removeListener(listener, device);
|
||||
}
|
||||
// TODO if tracks are always associated to only one device mapping, we could remove tracks here
|
||||
// TODO remove listener from store if there is no mapping left for the DeviceIdentifier
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void mappingAdded(DeviceMappingWithRegattaLogEvent<WithID> mapping) {
|
||||
// The listener is first added to not lose any fix after loading the initial fixes and adding the listener.
|
||||
sensorFixStore.addListener(listener, mapping.getDevice());
|
||||
loadFixes(getTrackingTimeRange().intersection(mapping.getTimeRange()), mapping);
|
||||
}
|
||||
|
||||
|
||||
+96
-90
@@ -30,6 +30,8 @@ import com.sap.sailing.domain.tracking.DynamicTrack;
|
||||
import com.sap.sse.common.TimePoint;
|
||||
import com.sap.sse.common.Timed;
|
||||
import com.sap.sse.common.WithID;
|
||||
import com.sap.sse.concurrent.LockUtil;
|
||||
import com.sap.sse.concurrent.NamedReentrantReadWriteLock;
|
||||
|
||||
/**
|
||||
* Holds DeviceMappings to make it possible to track changes. This makes it possible to only process mappings that
|
||||
@@ -46,6 +48,16 @@ public abstract class RegattaLogDeviceMappings<ItemT extends WithID> {
|
||||
* sure to not produce a memory leak.
|
||||
*/
|
||||
private final Set<RegattaLog> knownRegattaLogs = new HashSet<>();
|
||||
|
||||
/**
|
||||
* Lock object to be used when accessing {@link #knownRegattaLogs}.
|
||||
*/
|
||||
private final NamedReentrantReadWriteLock knownRegattaLogsLock;
|
||||
|
||||
/**
|
||||
* Lock object to be used when accessing {@link #mappings} or {@link #mappingsByDevice}.
|
||||
*/
|
||||
private final NamedReentrantReadWriteLock mappingsLock;
|
||||
|
||||
private final Map<ItemT, List<DeviceMappingWithRegattaLogEvent<ItemT>>> mappings = new HashMap<>();
|
||||
private final Map<DeviceIdentifier, List<DeviceMappingWithRegattaLogEvent<ItemT>>> mappingsByDevice = new HashMap<>();
|
||||
@@ -84,11 +96,16 @@ public abstract class RegattaLogDeviceMappings<ItemT extends WithID> {
|
||||
};
|
||||
};
|
||||
|
||||
public RegattaLogDeviceMappings(Iterable<RegattaLog> initialRegattaLogs) {
|
||||
public RegattaLogDeviceMappings(Iterable<RegattaLog> initialRegattaLogs, String raceNameForLock) {
|
||||
mappingsLock = new NamedReentrantReadWriteLock("DeviceMapping lock for race " + raceNameForLock, false);
|
||||
knownRegattaLogsLock = new NamedReentrantReadWriteLock("Lock for known RegattaLogs of race " + raceNameForLock, false);
|
||||
final boolean hasRegattaLogs;
|
||||
synchronized (knownRegattaLogs) {
|
||||
LockUtil.lockForWrite(knownRegattaLogsLock);
|
||||
try {
|
||||
initialRegattaLogs.forEach(this::addRegattaLogUnlocked);
|
||||
hasRegattaLogs = !knownRegattaLogs.isEmpty();
|
||||
} finally {
|
||||
LockUtil.unlockAfterWrite(knownRegattaLogsLock);
|
||||
}
|
||||
if (hasRegattaLogs) {
|
||||
updateMappings();
|
||||
@@ -96,9 +113,7 @@ public abstract class RegattaLogDeviceMappings<ItemT extends WithID> {
|
||||
}
|
||||
|
||||
public void addRegattaLog(RegattaLog regattaLog) {
|
||||
synchronized (knownRegattaLogs) {
|
||||
addRegattaLogUnlocked(regattaLog);
|
||||
}
|
||||
LockUtil.executeWithWriteLock(knownRegattaLogsLock, () -> addRegattaLogUnlocked(regattaLog));
|
||||
updateMappings();
|
||||
}
|
||||
|
||||
@@ -108,39 +123,19 @@ public abstract class RegattaLogDeviceMappings<ItemT extends WithID> {
|
||||
}
|
||||
|
||||
public void stop() {
|
||||
synchronized (knownRegattaLogs) {
|
||||
LockUtil.executeWithWriteLock(knownRegattaLogsLock, () -> {
|
||||
knownRegattaLogs.forEach((log) -> log.removeListener(regattaLogEventVisitor));
|
||||
knownRegattaLogs.clear();
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
protected void updateMappings() {
|
||||
try {
|
||||
updateMappings(true);
|
||||
} catch (Exception e) {
|
||||
logger.warning("Could not update device mappings");
|
||||
logger.log(Level.WARNING, "Could not update device mappings", e);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Used internally to access {@link #mappings} with a defensive copy to not run into concurrency issues if this is
|
||||
* called while an update is done.
|
||||
*
|
||||
* @return the mappings
|
||||
*/
|
||||
private synchronized Map<ItemT, List<DeviceMappingWithRegattaLogEvent<ItemT>>> getDeviceMappings() {
|
||||
return new HashMap<>(mappings);
|
||||
}
|
||||
|
||||
/**
|
||||
* Used internally to access {@link #mappingsByDevice} with a defensive copy to not run into concurrency issues if
|
||||
* this is called while an update is done.
|
||||
*
|
||||
* @return the mappings
|
||||
*/
|
||||
private synchronized Map<DeviceIdentifier, List<DeviceMappingWithRegattaLogEvent<ItemT>>> getMappingsByDevice() {
|
||||
return new HashMap<>(mappingsByDevice);
|
||||
}
|
||||
|
||||
/**
|
||||
* Calls the given callback for every known mapping that's currently known.
|
||||
@@ -157,7 +152,7 @@ public abstract class RegattaLogDeviceMappings<ItemT extends WithID> {
|
||||
* @param callback the callback to call for every mapped device
|
||||
*/
|
||||
public void forEachDevice(Consumer<DeviceIdentifier> callback) {
|
||||
getMappingsByDevice().keySet().forEach(callback::accept);
|
||||
LockUtil.executeWithReadLock(mappingsLock, () -> mappingsByDevice.keySet().forEach(callback::accept));
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -166,12 +161,14 @@ public abstract class RegattaLogDeviceMappings<ItemT extends WithID> {
|
||||
* @param callback the callback to call for every known mapping
|
||||
*/
|
||||
public void forEachMapping(BiConsumer<ItemT, DeviceMappingWithRegattaLogEvent<ItemT>> callback) {
|
||||
for (Map.Entry<ItemT, List<DeviceMappingWithRegattaLogEvent<ItemT>>> entry : getDeviceMappings().entrySet()) {
|
||||
ItemT item = entry.getKey();
|
||||
for (DeviceMappingWithRegattaLogEvent<ItemT> mapping : entry.getValue()) {
|
||||
callback.accept(item, mapping);
|
||||
LockUtil.executeWithReadLock(mappingsLock, () -> {
|
||||
for (Map.Entry<ItemT, List<DeviceMappingWithRegattaLogEvent<ItemT>>> entry : mappings.entrySet()) {
|
||||
ItemT item = entry.getKey();
|
||||
for (DeviceMappingWithRegattaLogEvent<ItemT> mapping : entry.getValue()) {
|
||||
callback.accept(item, mapping);
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -188,16 +185,27 @@ public abstract class RegattaLogDeviceMappings<ItemT extends WithID> {
|
||||
*/
|
||||
public void forEachMappingOfDeviceIncludingTimePoint(DeviceIdentifier device, TimePoint timePoint,
|
||||
Consumer<DeviceMappingWithRegattaLogEvent<ItemT>> callback) {
|
||||
List<DeviceMappingWithRegattaLogEvent<ItemT>> mappingsForDevice;
|
||||
synchronized (this) {
|
||||
mappingsForDevice = mappingsByDevice.get(device);
|
||||
}
|
||||
if (mappingsForDevice != null) {
|
||||
for (DeviceMappingWithRegattaLogEvent<ItemT> mapping : mappingsForDevice) {
|
||||
if (mapping.getTimeRange().includes(timePoint)) {
|
||||
callback.accept(mapping);
|
||||
LockUtil.executeWithReadLock(mappingsLock, () -> {
|
||||
List<DeviceMappingWithRegattaLogEvent<ItemT>> mappingsForDevice = mappingsByDevice.get(device);
|
||||
if (mappingsForDevice != null) {
|
||||
for (DeviceMappingWithRegattaLogEvent<ItemT> mapping : mappingsForDevice) {
|
||||
if (mapping.getTimeRange().includes(timePoint)) {
|
||||
callback.accept(mapping);
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* @return true if there is at least one mapping for the given {@link DeviceIdentifier}, false otherwise
|
||||
*/
|
||||
public boolean hasMappingForDevice(DeviceIdentifier device) {
|
||||
LockUtil.lockForRead(mappingsLock);
|
||||
try {
|
||||
return mappingsByDevice.containsKey(device);
|
||||
} finally {
|
||||
LockUtil.unlockAfterRead(mappingsLock);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -214,9 +222,7 @@ public abstract class RegattaLogDeviceMappings<ItemT extends WithID> {
|
||||
}
|
||||
|
||||
protected void forEachRegattaLog(Consumer<RegattaLog> regattaLogConsumer) {
|
||||
synchronized (knownRegattaLogs) {
|
||||
knownRegattaLogs.forEach(regattaLogConsumer);
|
||||
}
|
||||
LockUtil.executeWithReadLock(knownRegattaLogsLock, () -> knownRegattaLogs.forEach(regattaLogConsumer));
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -233,58 +239,58 @@ public abstract class RegattaLogDeviceMappings<ItemT extends WithID> {
|
||||
private final <FixT extends Timed, TrackT extends DynamicTrack<FixT>> void updateMappings(boolean loadIfNotCovered) {
|
||||
// TODO remove fixes, if mappings have been removed
|
||||
// check if there are new time ranges not covered so far
|
||||
Map<ItemT, List<DeviceMappingWithRegattaLogEvent<ItemT>>> newMappings = calculateMappings();
|
||||
updateMappings(newMappings, loadIfNotCovered);
|
||||
}
|
||||
|
||||
/**
|
||||
* Used internally to update the internal set of mappings to the state of the new mappings.
|
||||
* This ensures that the internal state is updated and all new and changed mappings are correctly processed.
|
||||
*
|
||||
* @param newMappings the new mappings
|
||||
* @param loadIfNotCovered true to inform about new/changed/removed mappings
|
||||
*/
|
||||
private synchronized void updateMappings(Map<ItemT, List<DeviceMappingWithRegattaLogEvent<ItemT>>> newMappings,
|
||||
boolean loadIfNotCovered) {
|
||||
if (loadIfNotCovered) {
|
||||
Set<ItemT> itemsToProcess = new HashSet<ItemT>(mappings.keySet());
|
||||
itemsToProcess.addAll(newMappings.keySet());
|
||||
for(ItemT item : itemsToProcess) {
|
||||
if(!newMappings.containsKey(item)) {
|
||||
mappings.get(item).forEach(this::mappingRemoved);
|
||||
} else {
|
||||
final List<DeviceMappingWithRegattaLogEvent<ItemT>> oldMappings = mappings.containsKey(item)
|
||||
? mappings.get(item) : Collections.emptyList();
|
||||
|
||||
for (DeviceMappingWithRegattaLogEvent<ItemT> newMapping : newMappings.get(item)) {
|
||||
DeviceMappingWithRegattaLogEvent<ItemT> oldMapping = findAndRemoveMapping(newMapping,
|
||||
oldMappings);
|
||||
if (oldMapping == null) {
|
||||
mappingAdded(newMapping);
|
||||
} else if (newMapping.getTimeRange().equals(oldMapping.getTimeRange())) {
|
||||
mappingChanged(oldMapping, newMapping);
|
||||
}
|
||||
final Map<ItemT, List<DeviceMappingWithRegattaLogEvent<ItemT>>> newMappings = calculateMappings();
|
||||
final Map<ItemT, List<DeviceMappingWithRegattaLogEvent<ItemT>>> oldMappings = new HashMap<>();
|
||||
LockUtil.lockForWrite(mappingsLock);
|
||||
try {
|
||||
oldMappings.putAll(mappings);
|
||||
|
||||
mappings.clear();
|
||||
mappings.putAll(newMappings);
|
||||
mappingsByDevice.clear();
|
||||
for (ItemT item : newMappings.keySet()) {
|
||||
for (DeviceMappingWithRegattaLogEvent<ItemT> mapping : newMappings.get(item)) {
|
||||
List<DeviceMappingWithRegattaLogEvent<ItemT>> list = mappingsByDevice.get(mapping.getDevice());
|
||||
if (list == null) {
|
||||
list = new ArrayList<>();
|
||||
mappingsByDevice.put(mapping.getDevice(), list);
|
||||
}
|
||||
oldMappings.forEach(this::mappingRemoved);
|
||||
list.add(mapping);
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
LockUtil.unlockAfterWrite(mappingsLock);
|
||||
}
|
||||
|
||||
mappings.clear();
|
||||
mappings.putAll(newMappings);
|
||||
mappingsByDevice.clear();
|
||||
for (ItemT item : newMappings.keySet()) {
|
||||
for (DeviceMappingWithRegattaLogEvent<ItemT> mapping : newMappings.get(item)) {
|
||||
List<DeviceMappingWithRegattaLogEvent<ItemT>> list = mappingsByDevice.get(mapping.getDevice());
|
||||
if (list == null) {
|
||||
list = new ArrayList<>();
|
||||
mappingsByDevice.put(mapping.getDevice(), list);
|
||||
}
|
||||
list.add(mapping);
|
||||
}
|
||||
if (loadIfNotCovered) {
|
||||
calculateDiff(oldMappings, newMappings);
|
||||
}
|
||||
}
|
||||
|
||||
private void calculateDiff(Map<ItemT, List<DeviceMappingWithRegattaLogEvent<ItemT>>> previousMappings,
|
||||
Map<ItemT, List<DeviceMappingWithRegattaLogEvent<ItemT>>> newMappings) {
|
||||
Set<ItemT> itemsToProcess = new HashSet<ItemT>(previousMappings.keySet());
|
||||
itemsToProcess.addAll(newMappings.keySet());
|
||||
for(ItemT item : itemsToProcess) {
|
||||
if(!newMappings.containsKey(item)) {
|
||||
previousMappings.get(item).forEach(this::mappingRemoved);
|
||||
} else {
|
||||
final List<DeviceMappingWithRegattaLogEvent<ItemT>> oldMappings = previousMappings.containsKey(item)
|
||||
? previousMappings.get(item) : Collections.emptyList();
|
||||
|
||||
for (DeviceMappingWithRegattaLogEvent<ItemT> newMapping : newMappings.get(item)) {
|
||||
DeviceMappingWithRegattaLogEvent<ItemT> oldMapping = findAndRemoveMapping(newMapping,
|
||||
oldMappings);
|
||||
if (oldMapping == null) {
|
||||
mappingAdded(newMapping);
|
||||
} else if (newMapping.getTimeRange().equals(oldMapping.getTimeRange())) {
|
||||
mappingChanged(oldMapping, newMapping);
|
||||
}
|
||||
}
|
||||
oldMappings.forEach(this::mappingRemoved);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Called when a {@link DeviceMapping} was removed.
|
||||
*
|
||||
|
||||
+4
@@ -20,6 +20,10 @@ public enum EmptySensorFixStore implements SensorFixStore {
|
||||
public void removeListener(FixReceivedListener<? extends Timed> listener) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void removeListener(FixReceivedListener<? extends Timed> listener, DeviceIdentifier device) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public TimeRange getTimeRangeCoveredByFixes(DeviceIdentifier device) {
|
||||
return null;
|
||||
|
||||
+5
@@ -59,6 +59,11 @@ public interface SensorFixStore {
|
||||
*/
|
||||
void removeListener(FixReceivedListener<? extends Timed> listener);
|
||||
|
||||
/**
|
||||
* Remove the registrations of the listener for the given device.
|
||||
*/
|
||||
void removeListener(FixReceivedListener<? extends Timed> listener, DeviceIdentifier device);
|
||||
|
||||
TimeRange getTimeRangeCoveredByFixes(DeviceIdentifier device) throws TransformationException,
|
||||
NoCorrespondingServiceRegisteredException;
|
||||
|
||||
|
||||
@@ -417,5 +417,30 @@ public class LockUtil {
|
||||
message.append(formatStackTrace(stackTrace));
|
||||
message.append('\n');
|
||||
}
|
||||
|
||||
|
||||
/**
|
||||
* Convenience method to execute a {@link Runnable} while the given {@link NamedReentrantReadWriteLock} is locked
|
||||
* for read. Ensures, that unlock is done in a finally block.
|
||||
*/
|
||||
public static void executeWithReadLock(NamedReentrantReadWriteLock lock, Runnable runnable) {
|
||||
lockForRead(lock);
|
||||
try {
|
||||
runnable.run();
|
||||
} finally {
|
||||
unlockAfterRead(lock);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Convenience method to execute a {@link Runnable} while the given {@link NamedReentrantReadWriteLock} is locked
|
||||
* for write. Ensures, that unlock is done in a finally block.
|
||||
*/
|
||||
public static void executeWithWriteLock(NamedReentrantReadWriteLock lock, Runnable runnable) {
|
||||
lockForWrite(lock);
|
||||
try {
|
||||
runnable.run();
|
||||
} finally {
|
||||
unlockAfterWrite(lock);
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user