share Igimi LiveDataConnection for equal combinations of device serial numbers

This commit is contained in:
Axel Uhl
2014-01-31 16:49:36 +01:00
parent a3779c9878
commit 454487fd2c
8 changed files with 129 additions and 7 deletions
@@ -158,7 +158,7 @@ public class WebSocketTest {
// the data from baur@stg-academy.org, particularly containing the Berlin test data
Account account = igtimiConnectionFactory.registerAccountForWhichClientIsAuthorized("9fded995cf21c8ed91ddaec13b220e8d5e44c65808d22ec2b1b7c32261121f26");
IgtimiConnection conn = igtimiConnectionFactory.connect(account);
LiveDataConnection liveDataConnection = conn.createLiveConnection(Collections.singleton("GA-EN-AAEJ"));
LiveDataConnection liveDataConnection = conn.getOrCreateLiveConnection(Collections.singleton("GA-EN-AAEJ"));
liveDataConnection.addListener(new BulkFixReceiver() {
@Override
public void received(Iterable<Fix> fixes) {
@@ -96,7 +96,7 @@ public interface IgtimiConnection {
*
* @return a connection that the caller can use to stop the live feed by calling {@link LiveDataConnection#stop()}.
*/
LiveDataConnection createLiveConnection(Iterable<String> deviceSerialNumbers) throws Exception;
LiveDataConnection getOrCreateLiveConnection(Iterable<String> deviceSerialNumbers) throws Exception;
/**
* @param sessionIds
@@ -35,7 +35,8 @@ import com.sap.sailing.domain.igtimiadapter.User;
import com.sap.sailing.domain.igtimiadapter.datatypes.Fix;
import com.sap.sailing.domain.igtimiadapter.datatypes.Type;
import com.sap.sailing.domain.igtimiadapter.shared.IgtimiWindReceiver;
import com.sap.sailing.domain.igtimiadapter.websocket.WebSocketConnectionManager;
import com.sap.sailing.domain.igtimiadapter.websocket.LiveDataConnectionFactory;
import com.sap.sailing.domain.igtimiadapter.websocket.LiveDataConnectionFactoryImpl;
import com.sap.sailing.domain.tracking.DynamicTrack;
import com.sap.sailing.domain.tracking.DynamicTrackedRace;
import com.sap.sailing.domain.tracking.TrackedRace;
@@ -45,10 +46,12 @@ public class IgtimiConnectionImpl implements IgtimiConnection {
private static final Logger logger = Logger.getLogger(IgtimiConnectionImpl.class.getName());
private final Account account;
private final IgtimiConnectionFactoryImpl connectionFactory;
private final LiveDataConnectionFactory liveDataConnectionFactory;
public IgtimiConnectionImpl(IgtimiConnectionFactoryImpl connectionFactory, Account account) {
this.connectionFactory = connectionFactory;
this.account = account;
liveDataConnectionFactory = new LiveDataConnectionFactoryImpl(connectionFactory, account);
}
@Override
@@ -184,8 +187,8 @@ public class IgtimiConnectionImpl implements IgtimiConnection {
@Override
public LiveDataConnection createLiveConnection(Iterable<String> deviceSerialNumbers) throws Exception {
return new WebSocketConnectionManager(connectionFactory, deviceSerialNumbers, getAccount());
public LiveDataConnection getOrCreateLiveConnection(Iterable<String> deviceSerialNumbers) throws Exception {
return liveDataConnectionFactory.getOrCreateLiveDataConnection(deviceSerialNumbers);
}
private DynamicTrack<Fix> getOrCreateTrack(Map<String, Map<Type, DynamicTrack<Fix>>> result,
@@ -46,7 +46,7 @@ public class IgtimiWindTracker extends AbstractWindTracker implements WindTracke
IgtimiConnection connection = connectionFactory.connect(account);
Iterable<String> devicesWeShouldListenTo = connection.getWindDevices();
if (!stopping) {
LiveDataConnection liveConnection = connection.createLiveConnection(devicesWeShouldListenTo);
LiveDataConnection liveConnection = connection.getOrCreateLiveConnection(devicesWeShouldListenTo);
IgtimiWindReceiver windReceiver = new IgtimiWindReceiver(devicesWeShouldListenTo);
liveConnection.addListener(windReceiver);
windReceiver.addListener(new WindListenerSendingToTrackedRace(Collections.singleton(getTrackedRace()), windTrackerFactory));
@@ -0,0 +1,16 @@
package com.sap.sailing.domain.igtimiadapter.websocket;
import com.sap.sailing.domain.igtimiadapter.LiveDataConnection;
/**
* Helps bundling live data connections for the same set of devices. Clients can request a connection for a set of devices.
* If one already exists, a wrapper to it is returned. This wrapper's {@link LiveDataConnection#stop()} method work such that
* it only decrements a usage counter in this factory (and does so at most once), such that the actual connection is only
* terminated if the last client has stopped using it.
*
* @author Axel Uhl (D043530)
*
*/
public interface LiveDataConnectionFactory {
LiveDataConnection getOrCreateLiveDataConnection(Iterable<String> deviceSerialNumbers) throws Exception;
}
@@ -0,0 +1,66 @@
package com.sap.sailing.domain.igtimiadapter.websocket;
import java.util.HashMap;
import java.util.HashSet;
import java.util.Map;
import java.util.Set;
import java.util.logging.Logger;
import com.sap.sailing.domain.common.impl.Util;
import com.sap.sailing.domain.igtimiadapter.Account;
import com.sap.sailing.domain.igtimiadapter.LiveDataConnection;
import com.sap.sailing.domain.igtimiadapter.impl.IgtimiConnectionFactoryImpl;
public class LiveDataConnectionFactoryImpl implements LiveDataConnectionFactory {
private static final Logger logger = Logger.getLogger(LiveDataConnectionFactoryImpl.class.getName());
private final IgtimiConnectionFactoryImpl connectionFactory;
private final Account account;
private final Map<Set<String>, LiveDataConnection> dataConnectionsForDeviceSerialNumbers;
private final Map<LiveDataConnection, Set<String>> deviceSerialNumersForDataConnections;
private final Map<LiveDataConnection, Integer> usageCounts;
public LiveDataConnectionFactoryImpl(IgtimiConnectionFactoryImpl connectionFactory, Account account) {
this.connectionFactory = connectionFactory;
this.account = account;
dataConnectionsForDeviceSerialNumbers = new HashMap<>();
deviceSerialNumersForDataConnections = new HashMap<>();
usageCounts = new HashMap<>();
}
@Override
public synchronized LiveDataConnection getOrCreateLiveDataConnection(Iterable<String> deviceSerialNumbers) throws Exception {
Set<String> deviceSerialNumbersAsSet = new HashSet<>();
Util.addAll(deviceSerialNumbers, deviceSerialNumbersAsSet);
LiveDataConnection result = dataConnectionsForDeviceSerialNumbers.get(dataConnectionsForDeviceSerialNumbers);
if (result == null) {
result = new WebSocketConnectionManager(connectionFactory, deviceSerialNumbers, account);
dataConnectionsForDeviceSerialNumbers.put(deviceSerialNumbersAsSet, result);
deviceSerialNumersForDataConnections.put(result, deviceSerialNumbersAsSet);
}
Integer usageCount = usageCounts.get(result);
if (usageCount == null) {
usageCount = 0;
}
usageCount++;
usageCounts.put(result, usageCount);
return result;
}
public synchronized void stop(LiveDataConnection actualConnection) throws Exception {
Integer usageCount = usageCounts.get(actualConnection);
if (usageCount == null || usageCount == 0) {
logger.warning("Strange: the Igtimi live data connection "+actualConnection+" is released by another client although no client should be using it anymore.");
} else {
usageCount--;
if (usageCount == 0) {
usageCounts.remove(actualConnection);
Set<String> deviceSerialNumbersAsSet = deviceSerialNumersForDataConnections.remove(actualConnection);
dataConnectionsForDeviceSerialNumbers.remove(deviceSerialNumbersAsSet);
actualConnection.stop();
} else {
usageCounts.put(actualConnection, usageCount);
}
}
}
}
@@ -0,0 +1,37 @@
package com.sap.sailing.domain.igtimiadapter.websocket;
import com.sap.sailing.domain.igtimiadapter.BulkFixReceiver;
import com.sap.sailing.domain.igtimiadapter.LiveDataConnection;
public class LiveDataConnectionWrapper implements LiveDataConnection {
private final LiveDataConnectionFactoryImpl factory;
private final LiveDataConnection actualConnection;
private boolean stopCalled;
protected LiveDataConnectionWrapper(LiveDataConnectionFactoryImpl factory, LiveDataConnection actualConnection) {
super();
this.factory = factory;
this.actualConnection = actualConnection;
}
@Override
public synchronized void stop() throws Exception {
if (!stopCalled) {
factory.stop(actualConnection);
stopCalled = true;
}
}
@Override
public boolean waitForConnection(long timeoutInMillis) throws InterruptedException {
return actualConnection.waitForConnection(timeoutInMillis);
}
@Override
public void addListener(BulkFixReceiver listener) {
actualConnection.addListener(listener);
}
}
@@ -176,7 +176,7 @@ public class WindStatusServlet extends SailingServerHttpServlet {
if (account.getUser() != null) {
IgtimiConnection igtimiConnection = igtimiConnectionFactory.connect(account);
try {
LiveDataConnection liveDataConnection = igtimiConnection.createLiveConnection(igtimiConnection.getWindDevices());
LiveDataConnection liveDataConnection = igtimiConnection.getOrCreateLiveConnection(igtimiConnection.getWindDevices());
result = true;
liveDataConnection.addListener(new BulkFixReceiver() {
@Override