mirror of
https://github.com/eclipse-sailing-analytics/sailing-analytics.git
synced 2026-09-30 09:26:44 +00:00
Merge branch 'master' into 'gwt25'
Conflicts: java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/leaderboard/SortedCellTableWithStylableHeaders.java
This commit is contained in:
@@ -1,7 +1,7 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<classpath>
|
||||
<classpathentry kind="con" path="org.eclipse.pde.core.requiredPlugins"/>
|
||||
<classpathentry kind="src" path="src"/>
|
||||
<classpathentry kind="con" path="org.eclipse.jdt.launching.JRE_CONTAINER/org.eclipse.jdt.internal.debug.ui.launcher.StandardVMType/JavaSE-1.7"/>
|
||||
<classpathentry kind="output" path="bin"/>
|
||||
</classpath>
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<classpath>
|
||||
<classpathentry kind="con" path="org.eclipse.pde.core.requiredPlugins"/>
|
||||
<classpathentry kind="src" path="src"/>
|
||||
<classpathentry kind="con" path="org.eclipse.jdt.launching.JRE_CONTAINER/org.eclipse.jdt.internal.debug.ui.launcher.StandardVMType/JavaSE-1.7"/>
|
||||
<classpathentry kind="output" path="bin"/>
|
||||
</classpath>
|
||||
|
||||
@@ -1,28 +1,28 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<projectDescription>
|
||||
<name>com.sap.sailing.domain.swisstimingadapter.persistence</name>
|
||||
<comment></comment>
|
||||
<projects>
|
||||
</projects>
|
||||
<buildSpec>
|
||||
<buildCommand>
|
||||
<name>org.eclipse.jdt.core.javabuilder</name>
|
||||
<arguments>
|
||||
</arguments>
|
||||
</buildCommand>
|
||||
<buildCommand>
|
||||
<name>org.eclipse.pde.ManifestBuilder</name>
|
||||
<arguments>
|
||||
</arguments>
|
||||
</buildCommand>
|
||||
<buildCommand>
|
||||
<name>org.eclipse.pde.SchemaBuilder</name>
|
||||
<arguments>
|
||||
</arguments>
|
||||
</buildCommand>
|
||||
</buildSpec>
|
||||
<natures>
|
||||
<nature>org.eclipse.pde.PluginNature</nature>
|
||||
<nature>org.eclipse.jdt.core.javanature</nature>
|
||||
</natures>
|
||||
</projectDescription>
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<projectDescription>
|
||||
<name>com.sap.sailing.domain.swisstimingadapter.persistence</name>
|
||||
<comment></comment>
|
||||
<projects>
|
||||
</projects>
|
||||
<buildSpec>
|
||||
<buildCommand>
|
||||
<name>org.eclipse.jdt.core.javabuilder</name>
|
||||
<arguments>
|
||||
</arguments>
|
||||
</buildCommand>
|
||||
<buildCommand>
|
||||
<name>org.eclipse.pde.ManifestBuilder</name>
|
||||
<arguments>
|
||||
</arguments>
|
||||
</buildCommand>
|
||||
<buildCommand>
|
||||
<name>org.eclipse.pde.SchemaBuilder</name>
|
||||
<arguments>
|
||||
</arguments>
|
||||
</buildCommand>
|
||||
</buildSpec>
|
||||
<natures>
|
||||
<nature>org.eclipse.pde.PluginNature</nature>
|
||||
<nature>org.eclipse.jdt.core.javanature</nature>
|
||||
</natures>
|
||||
</projectDescription>
|
||||
|
||||
+8
-8
@@ -1,8 +1,8 @@
|
||||
#Thu Jan 19 09:52:00 CET 2012
|
||||
eclipse.preferences.version=1
|
||||
org.eclipse.jdt.core.compiler.codegen.inlineJsrBytecode=enabled
|
||||
org.eclipse.jdt.core.compiler.codegen.targetPlatform=1.7
|
||||
org.eclipse.jdt.core.compiler.compliance=1.7
|
||||
org.eclipse.jdt.core.compiler.problem.assertIdentifier=error
|
||||
org.eclipse.jdt.core.compiler.problem.enumIdentifier=error
|
||||
org.eclipse.jdt.core.compiler.source=1.7
|
||||
#Thu Jan 19 09:52:00 CET 2012
|
||||
eclipse.preferences.version=1
|
||||
org.eclipse.jdt.core.compiler.codegen.inlineJsrBytecode=enabled
|
||||
org.eclipse.jdt.core.compiler.codegen.targetPlatform=1.7
|
||||
org.eclipse.jdt.core.compiler.compliance=1.7
|
||||
org.eclipse.jdt.core.compiler.problem.assertIdentifier=error
|
||||
org.eclipse.jdt.core.compiler.problem.enumIdentifier=error
|
||||
org.eclipse.jdt.core.compiler.source=1.7
|
||||
|
||||
+4
-4
@@ -1,4 +1,4 @@
|
||||
#Tue Nov 08 17:33:37 CET 2011
|
||||
eclipse.preferences.version=1
|
||||
pluginProject.extensions=false
|
||||
resolve.requirebundle=false
|
||||
#Tue Nov 08 17:33:37 CET 2011
|
||||
eclipse.preferences.version=1
|
||||
pluginProject.extensions=false
|
||||
resolve.requirebundle=false
|
||||
|
||||
@@ -1,12 +1,13 @@
|
||||
<?xml version="1.0" encoding="UTF-8" standalone="no"?>
|
||||
<launchConfiguration type="org.eclipse.jdt.launching.localJavaApplication">
|
||||
<listAttribute key="org.eclipse.debug.core.MAPPED_RESOURCE_PATHS">
|
||||
<listEntry value="/com.sap.sailing.domain.swisstimingadapter.persistence/src/com/sap/sailing/domain/swisstimingadapter/persistence/StoreAndForward.java"/>
|
||||
</listAttribute>
|
||||
<listAttribute key="org.eclipse.debug.core.MAPPED_RESOURCE_TYPES">
|
||||
<listEntry value="1"/>
|
||||
</listAttribute>
|
||||
<stringAttribute key="org.eclipse.jdt.launching.MAIN_TYPE" value="com.sap.sailing.domain.swisstimingadapter.persistence.StoreAndForward"/>
|
||||
<stringAttribute key="org.eclipse.jdt.launching.PROGRAM_ARGUMENTS" value="3500 3501"/>
|
||||
<stringAttribute key="org.eclipse.jdt.launching.PROJECT_ATTR" value="com.sap.sailing.domain.swisstimingadapter.persistence"/>
|
||||
</launchConfiguration>
|
||||
<?xml version="1.0" encoding="UTF-8" standalone="no"?>
|
||||
<launchConfiguration type="org.eclipse.jdt.launching.localJavaApplication">
|
||||
<listAttribute key="org.eclipse.debug.core.MAPPED_RESOURCE_PATHS">
|
||||
<listEntry value="/com.sap.sailing.domain.swisstimingadapter.persistence/src/com/sap/sailing/domain/swisstimingadapter/persistence/StoreAndForward.java"/>
|
||||
</listAttribute>
|
||||
<listAttribute key="org.eclipse.debug.core.MAPPED_RESOURCE_TYPES">
|
||||
<listEntry value="1"/>
|
||||
</listAttribute>
|
||||
<stringAttribute key="org.eclipse.jdt.launching.MAIN_TYPE" value="com.sap.sailing.domain.swisstimingadapter.persistence.StoreAndForward"/>
|
||||
<stringAttribute key="org.eclipse.jdt.launching.PROGRAM_ARGUMENTS" value="3500 3501"/>
|
||||
<stringAttribute key="org.eclipse.jdt.launching.PROJECT_ATTR" value="com.sap.sailing.domain.swisstimingadapter.persistence"/>
|
||||
<stringAttribute key="org.eclipse.jdt.launching.VM_ARGUMENTS" value="-Djava.util.logging.config.file=${project_loc:com.sap.sailing.server}/../target/configuration/logging_debug.properties"/>
|
||||
</launchConfiguration>
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
source.. = src/
|
||||
output.. = bin/
|
||||
bin.includes = META-INF/,\
|
||||
.
|
||||
source.. = src/
|
||||
output.. = bin/
|
||||
bin.includes = META-INF/,\
|
||||
.
|
||||
|
||||
+341
-314
@@ -1,314 +1,341 @@
|
||||
package com.sap.sailing.domain.swisstimingadapter.persistence;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.io.InputStream;
|
||||
import java.io.OutputStream;
|
||||
import java.net.ServerSocket;
|
||||
import java.net.Socket;
|
||||
import java.net.UnknownHostException;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.logging.Logger;
|
||||
|
||||
import com.mongodb.BasicDBObject;
|
||||
import com.mongodb.DB;
|
||||
import com.mongodb.DBCollection;
|
||||
import com.mongodb.DBObject;
|
||||
import com.sap.sailing.domain.common.impl.Util.Pair;
|
||||
import com.sap.sailing.domain.swisstimingadapter.SailMasterMessage;
|
||||
import com.sap.sailing.domain.swisstimingadapter.SailMasterTransceiver;
|
||||
import com.sap.sailing.domain.swisstimingadapter.SwissTimingFactory;
|
||||
import com.sap.sailing.domain.swisstimingadapter.persistence.impl.CollectionNames;
|
||||
import com.sap.sailing.domain.swisstimingadapter.persistence.impl.FieldNames;
|
||||
import com.sap.sailing.mongodb.MongoDBConfiguration;
|
||||
import com.sap.sailing.mongodb.MongoDBService;
|
||||
|
||||
/**
|
||||
* Receives events from a SwissTiming SailMaster server, stores valid messages received persistently and forward them
|
||||
* to a port specified. The messages forwarded are augmented by sending a counter in ASCII encoding before the
|
||||
* message's <code>STX</code> start byte. This allows a receiver to optionally record the counter value after
|
||||
* having processed the message. When messages have to be retrieved from the database at a later point, a
|
||||
* client can request only those message starting at a specific counter value.<p>
|
||||
*
|
||||
* The connectivity to the SailMaster system can operate in one of two modes. Either a TCP connection to the
|
||||
* SailMaster is initiated by this client. When the connection is lost or a connect is unsuccessful, the client
|
||||
* will try again periodically. In the other mode of operation, this client will act as a TCP server and accepts
|
||||
* inbound requests from the SailMaster system (typically a bridge that forwards messages from multiple
|
||||
* SailMaster systems). In this mode of operation it's up to the SailMaster environment to re-initiate
|
||||
* connects after a connection loss.
|
||||
*
|
||||
* @author Axel Uhl (d043530)
|
||||
*
|
||||
*/
|
||||
public class StoreAndForward implements Runnable {
|
||||
private static final Logger logger = Logger.getLogger(StoreAndForward.class.getName());
|
||||
|
||||
private final DB db;
|
||||
private final int listenPort;
|
||||
private final SailMasterTransceiver transceiver;
|
||||
private final int portForClients;
|
||||
private long lastMessageCount;
|
||||
private boolean stopped;
|
||||
private final Thread clientListener;
|
||||
private final List<Socket> socketsToForwardTo;
|
||||
private final List<OutputStream> streamsToForwardTo;
|
||||
private boolean listeningForClients;
|
||||
private boolean receivingFromSailMaster;
|
||||
|
||||
private final DBCollection lastMessageCountCollection;
|
||||
|
||||
private final SwissTimingAdapterPersistence swissTimingAdapterPersistence;
|
||||
|
||||
private final SwissTimingFactory swissTimingFactory;
|
||||
|
||||
private final Thread storeAndForwardThread;
|
||||
|
||||
private Socket socket;
|
||||
|
||||
/**
|
||||
* Use of this server socket is optional and happens if and only if this object is operated in
|
||||
* "listening" mode. This means that the SailMaster system / bridge is expected to initiate TCP
|
||||
* connections to this object.
|
||||
*/
|
||||
private ServerSocket serverSocketListeningForSailMasterBridge;
|
||||
|
||||
private final int sailMasterPort;
|
||||
|
||||
private final String sailMasterHostname;
|
||||
|
||||
/**
|
||||
* When initialized using this constructor, the resulting object proactively connects and re-connects to the
|
||||
* SailMaster server specified by <code>sailMasterHostname</code>/<code>sailMasterPort</code>. The
|
||||
* {@link #serverSocketListeningForSailMasterBridge} remains <code>null</code> and {@link #listenPort} is set
|
||||
* to <code>-1</code>, indicating that this object is not listening on any port for incoming
|
||||
* SailMaster connections.
|
||||
*/
|
||||
public StoreAndForward(String sailMasterHostname, int sailMasterPort, int portForClients, SwissTimingFactory swissTimingFactory,
|
||||
SwissTimingAdapterPersistence swissTimingAdapterPersistence, MongoDBService mongoDBService) throws InterruptedException, IOException {
|
||||
this.db = mongoDBService.getDB();
|
||||
this.listenPort = -1;
|
||||
this.transceiver = swissTimingFactory.createSailMasterTransceiver();
|
||||
this.portForClients = portForClients;
|
||||
this.sailMasterHostname = sailMasterHostname;
|
||||
this.sailMasterPort = sailMasterPort;
|
||||
this.streamsToForwardTo = new ArrayList<OutputStream>();
|
||||
this.socketsToForwardTo = new ArrayList<Socket>();
|
||||
this.swissTimingAdapterPersistence = swissTimingAdapterPersistence;
|
||||
this.swissTimingFactory = swissTimingFactory;
|
||||
lastMessageCountCollection = db.getCollection(CollectionNames.LAST_MESSAGE_COUNT.name());
|
||||
lastMessageCount = getLastMessageCount();
|
||||
clientListener = createClientListenerThread(portForClients);
|
||||
clientListener.start();
|
||||
storeAndForwardThread = new Thread(this, "StoreAndForward");
|
||||
storeAndForwardThread.start();
|
||||
synchronized (this) {
|
||||
while (!listeningForClients) {
|
||||
wait();
|
||||
}
|
||||
while (!receivingFromSailMaster) {
|
||||
wait();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates a storing message forwarder in listening mode. In this mode, this object won't actively try to open
|
||||
* TCP connections to a SailMaster system / bridge but instead listen for inbound TCP connections on port
|
||||
* <code>listenPort</code>.
|
||||
*
|
||||
* @param listenPort
|
||||
* listens on this port for messages coming in from a real SwissTiming SailMaster
|
||||
* @param portForClients
|
||||
* clients can connect to this port and will receive forwarded and sequence-numbered messages over those
|
||||
* sockets
|
||||
* @throws IOException
|
||||
*/
|
||||
public StoreAndForward(final int listenPort, final int portForClients, SwissTimingFactory swissTimingFactory,
|
||||
SwissTimingAdapterPersistence swissTimingAdapterPersistence, MongoDBService mongoDBService) throws InterruptedException, IOException {
|
||||
this.db = mongoDBService.getDB();
|
||||
this.listenPort = listenPort;
|
||||
this.transceiver = swissTimingFactory.createSailMasterTransceiver();
|
||||
this.portForClients = portForClients;
|
||||
this.streamsToForwardTo = new ArrayList<OutputStream>();
|
||||
this.socketsToForwardTo = new ArrayList<Socket>();
|
||||
this.swissTimingAdapterPersistence = swissTimingAdapterPersistence;
|
||||
this.swissTimingFactory = swissTimingFactory;
|
||||
this.sailMasterHostname = null;
|
||||
this.sailMasterPort = -1;
|
||||
lastMessageCountCollection = db.getCollection(CollectionNames.LAST_MESSAGE_COUNT.name());
|
||||
lastMessageCount = getLastMessageCount();
|
||||
serverSocketListeningForSailMasterBridge = new ServerSocket(listenPort);
|
||||
clientListener = createClientListenerThread(portForClients);
|
||||
clientListener.start();
|
||||
storeAndForwardThread = new Thread(this, "StoreAndForward");
|
||||
storeAndForwardThread.start();
|
||||
synchronized (this) {
|
||||
while (!listeningForClients) {
|
||||
wait();
|
||||
}
|
||||
while (!receivingFromSailMaster) {
|
||||
wait();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private long getLastMessageCount() {
|
||||
DBObject lastMessageCountRecord = lastMessageCountCollection.findOne();
|
||||
if (lastMessageCountRecord == null) {
|
||||
lastMessageCountCollection.insert(new BasicDBObject().append(FieldNames.LAST_MESSAGE_COUNT.name(), 0l));
|
||||
}
|
||||
return lastMessageCountRecord == null ? 0 : (Long) lastMessageCountRecord.get(FieldNames.LAST_MESSAGE_COUNT.name());
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns <code>true</code> if and only if this object is listening for incoming SailMaster TCP
|
||||
* connections instead of actively connecting / reconnecting to a SailMaster server by itself.
|
||||
*/
|
||||
private boolean isInSailMasterListeningMode() {
|
||||
return serverSocketListeningForSailMasterBridge != null;
|
||||
}
|
||||
|
||||
private Thread createClientListenerThread(final int portForClients) {
|
||||
return new Thread(new Runnable() {
|
||||
public void run() {
|
||||
ServerSocket ss;
|
||||
try {
|
||||
synchronized (StoreAndForward.this) {
|
||||
ss = new ServerSocket(portForClients);
|
||||
listeningForClients = true;
|
||||
logger.info("StoreAndForward listening for clients on port "+portForClients);
|
||||
StoreAndForward.this.notifyAll();
|
||||
}
|
||||
while (!stopped) {
|
||||
Socket s = ss.accept();
|
||||
logger.info("StoreAndForward received connector's connect request on port "+portForClients);
|
||||
if (!stopped) {
|
||||
synchronized (StoreAndForward.this) {
|
||||
socketsToForwardTo.add(s);
|
||||
streamsToForwardTo.add(s.getOutputStream());
|
||||
}
|
||||
} else {
|
||||
s.close();
|
||||
}
|
||||
}
|
||||
ss.close();
|
||||
logger.info("StoreAndForward client listener thread stopped.");
|
||||
} catch (IOException e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
}
|
||||
}, "StoreAndForwardClientListener");
|
||||
}
|
||||
|
||||
/**
|
||||
* Depending on the mode of operation, accepts an inbound connect from a SailMaster server / bridge or
|
||||
* actively initiates a TCP connection to the SailMaster address/port configured. When this method
|
||||
* returns, the {@link #socket} holds an open and connected socket.
|
||||
* @throws IOException
|
||||
*/
|
||||
private void establishConnection() throws IOException {
|
||||
if (isInSailMasterListeningMode()) {
|
||||
synchronized (this) {
|
||||
receivingFromSailMaster = true;
|
||||
logger.info("StoreAndForward waiting for inbound SailMaster connections on port "+listenPort);
|
||||
notifyAll();
|
||||
}
|
||||
socket = serverSocketListeningForSailMasterBridge.accept();
|
||||
logger.info("StoreAndForward received SailMaster connect on port "+listenPort);
|
||||
} else {
|
||||
synchronized (this) {
|
||||
receivingFromSailMaster = true;
|
||||
logger.info("StoreAndForward issuing SailMaster connection to "+sailMasterHostname+":"+sailMasterPort);
|
||||
notifyAll();
|
||||
}
|
||||
socket = new Socket(sailMasterHostname, sailMasterPort);
|
||||
logger.info("StoreAndForward connections to SailMaster "+sailMasterHostname+":"+sailMasterPort+" established");
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Stops execution after having received the next message
|
||||
*/
|
||||
public void stop() throws UnknownHostException, IOException, InterruptedException {
|
||||
logger.entering(getClass().getName(), "stop");
|
||||
stopped = true;
|
||||
new Socket("localhost", portForClients); // this is to stop the client listener thread
|
||||
logger.info("joining clientListener thread "+clientListener);
|
||||
clientListener.join();
|
||||
socket.close(); // will let a read terminate abnormally
|
||||
logger.info("joining storeAndForwardThread "+storeAndForwardThread);
|
||||
storeAndForwardThread.join();
|
||||
}
|
||||
|
||||
public static void main(String[] args) throws InterruptedException, IOException {
|
||||
int listenPort = Integer.valueOf(args[0]);
|
||||
int clientPort = Integer.valueOf(args[1]);
|
||||
|
||||
MongoDBService mongoDBService = MongoDBService.INSTANCE;
|
||||
mongoDBService.setConfiguration(MongoDBConfiguration.getDefaultConfiguration());
|
||||
SwissTimingAdapterPersistence swissTimingAdapterPersistence = SwissTimingAdapterPersistence.INSTANCE;
|
||||
new StoreAndForward(listenPort, clientPort, SwissTimingFactory.INSTANCE, swissTimingAdapterPersistence, mongoDBService);
|
||||
}
|
||||
|
||||
public void run() {
|
||||
logger.entering(getClass().getName(), "run");
|
||||
try {
|
||||
while (!stopped) {
|
||||
try {
|
||||
establishConnection();
|
||||
InputStream is = socket.getInputStream();
|
||||
Pair<String, Long> messageAndOptionalSequenceNumber = transceiver.receiveMessage(is);
|
||||
// ignore any sequence number contained in the message; we'll create our own
|
||||
DBObject emptyQuery = new BasicDBObject();
|
||||
DBObject incrementLastMessageCountQuery = new BasicDBObject().
|
||||
append("$inc", new BasicDBObject().append(FieldNames.LAST_MESSAGE_COUNT.name(), 1));
|
||||
while (!stopped && messageAndOptionalSequenceNumber != null) {
|
||||
logger.fine("Received message: "+messageAndOptionalSequenceNumber.getA());
|
||||
DBObject newCountRecord = lastMessageCountCollection.findAndModify(emptyQuery, incrementLastMessageCountQuery);
|
||||
lastMessageCount = (Long) ((newCountRecord == null) ? 0 : newCountRecord.get(FieldNames.LAST_MESSAGE_COUNT.name()));
|
||||
SailMasterMessage message = swissTimingFactory.createMessage(messageAndOptionalSequenceNumber.getA(), lastMessageCount);
|
||||
swissTimingAdapterPersistence.storeSailMasterMessage(message);
|
||||
synchronized (this) {
|
||||
for (OutputStream os : streamsToForwardTo) {
|
||||
// write the sequence number of the message into the stream before actually writing the
|
||||
// SwissTiming message
|
||||
// TODO if forwarding to os doesn't work, e.g., because the socket was closed or the client died, remove os from streamsToForwardTo and the socket from socketsToForwardTo
|
||||
transceiver.sendMessage(message, os);
|
||||
}
|
||||
}
|
||||
if (!stopped) {
|
||||
messageAndOptionalSequenceNumber = transceiver.receiveMessage(is);
|
||||
}
|
||||
}
|
||||
for (OutputStream os : streamsToForwardTo) {
|
||||
os.close();
|
||||
}
|
||||
for (Socket socketToForwardTo : socketsToForwardTo) {
|
||||
socketToForwardTo.close();
|
||||
}
|
||||
} catch (Throwable e) {
|
||||
e.printStackTrace();
|
||||
if (!stopped) {
|
||||
logger.throwing(StoreAndForward.class.getName(), "Error during forwarding message. Continuing...", e);
|
||||
try {
|
||||
Thread.sleep(1000l); // wait a little bit before trying to re-establish a connection
|
||||
} catch (InterruptedException e1) {
|
||||
logger.throwing(StoreAndForward.class.getName(), "Can't find any sleep...", e1);
|
||||
}
|
||||
} else {
|
||||
logger.info("StoreAndForward socket was closed.");
|
||||
}
|
||||
}
|
||||
}
|
||||
logger.info("StoreAndForward is closing receiving server socket");
|
||||
if (serverSocketListeningForSailMasterBridge != null && !serverSocketListeningForSailMasterBridge.isClosed()) {
|
||||
serverSocketListeningForSailMasterBridge.close();
|
||||
}
|
||||
logger.info("Stopping StoreAndForward server.");
|
||||
} catch (IOException e) {
|
||||
logger.throwing(getClass().getName(), "run", e);
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
logger.exiting(getClass().getName(), "run");
|
||||
}
|
||||
}
|
||||
package com.sap.sailing.domain.swisstimingadapter.persistence;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.io.InputStream;
|
||||
import java.io.OutputStream;
|
||||
import java.net.ServerSocket;
|
||||
import java.net.Socket;
|
||||
import java.net.UnknownHostException;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.logging.Logger;
|
||||
|
||||
import com.mongodb.BasicDBObject;
|
||||
import com.mongodb.DB;
|
||||
import com.mongodb.DBCollection;
|
||||
import com.mongodb.DBObject;
|
||||
import com.sap.sailing.domain.common.impl.Util.Pair;
|
||||
import com.sap.sailing.domain.swisstimingadapter.SailMasterMessage;
|
||||
import com.sap.sailing.domain.swisstimingadapter.SailMasterTransceiver;
|
||||
import com.sap.sailing.domain.swisstimingadapter.SwissTimingFactory;
|
||||
import com.sap.sailing.domain.swisstimingadapter.persistence.impl.CollectionNames;
|
||||
import com.sap.sailing.domain.swisstimingadapter.persistence.impl.FieldNames;
|
||||
import com.sap.sailing.mongodb.MongoDBConfiguration;
|
||||
import com.sap.sailing.mongodb.MongoDBService;
|
||||
|
||||
/**
|
||||
* Receives events from a SwissTiming SailMaster server, stores valid messages received persistently and forward them
|
||||
* to a port specified. The messages forwarded are augmented by sending a counter in ASCII encoding before the
|
||||
* message's <code>STX</code> start byte. This allows a receiver to optionally record the counter value after
|
||||
* having processed the message. When messages have to be retrieved from the database at a later point, a
|
||||
* client can request only those message starting at a specific counter value.<p>
|
||||
*
|
||||
* The connectivity to the SailMaster system can operate in one of two modes. Either a TCP connection to the
|
||||
* SailMaster is initiated by this client. When the connection is lost or a connect is unsuccessful, the client
|
||||
* will try again periodically. In the other mode of operation, this client will act as a TCP server and accepts
|
||||
* inbound requests from the SailMaster system (typically a bridge that forwards messages from multiple
|
||||
* SailMaster systems). In this mode of operation it's up to the SailMaster environment to re-initiate
|
||||
* connects after a connection loss.
|
||||
*
|
||||
* @author Axel Uhl (d043530)
|
||||
*
|
||||
*/
|
||||
public class StoreAndForward implements Runnable {
|
||||
private static final Logger logger = Logger.getLogger(StoreAndForward.class.getName());
|
||||
|
||||
private final DB db;
|
||||
private final int listenPort;
|
||||
private final SailMasterTransceiver transceiver;
|
||||
private final int portForClients;
|
||||
private long lastMessageCount;
|
||||
private boolean stopped;
|
||||
private final Thread clientListener;
|
||||
private final List<Socket> socketsToForwardTo;
|
||||
private final List<OutputStream> streamsToForwardTo;
|
||||
private boolean listeningForClients;
|
||||
private boolean receivingFromSailMaster;
|
||||
|
||||
private final DBCollection lastMessageCountCollection;
|
||||
|
||||
private final SwissTimingAdapterPersistence swissTimingAdapterPersistence;
|
||||
|
||||
private final SwissTimingFactory swissTimingFactory;
|
||||
|
||||
private final Thread storeAndForwardThread;
|
||||
|
||||
private Socket socket;
|
||||
|
||||
/**
|
||||
* Use of this server socket is optional and happens if and only if this object is operated in
|
||||
* "listening" mode. This means that the SailMaster system / bridge is expected to initiate TCP
|
||||
* connections to this object.
|
||||
*/
|
||||
private ServerSocket serverSocketListeningForSailMasterBridge;
|
||||
|
||||
private final int sailMasterPort;
|
||||
|
||||
private final String sailMasterHostname;
|
||||
|
||||
/**
|
||||
* When initialized using this constructor, the resulting object proactively connects and re-connects to the
|
||||
* SailMaster server specified by <code>sailMasterHostname</code>/<code>sailMasterPort</code>. The
|
||||
* {@link #serverSocketListeningForSailMasterBridge} remains <code>null</code> and {@link #listenPort} is set
|
||||
* to <code>-1</code>, indicating that this object is not listening on any port for incoming
|
||||
* SailMaster connections.
|
||||
*/
|
||||
public StoreAndForward(String sailMasterHostname, int sailMasterPort, int portForClients, SwissTimingFactory swissTimingFactory,
|
||||
SwissTimingAdapterPersistence swissTimingAdapterPersistence, MongoDBService mongoDBService) throws InterruptedException, IOException {
|
||||
this.db = mongoDBService.getDB();
|
||||
this.listenPort = -1;
|
||||
this.transceiver = swissTimingFactory.createSailMasterTransceiver();
|
||||
this.portForClients = portForClients;
|
||||
this.sailMasterHostname = sailMasterHostname;
|
||||
this.sailMasterPort = sailMasterPort;
|
||||
this.streamsToForwardTo = new ArrayList<OutputStream>();
|
||||
this.socketsToForwardTo = new ArrayList<Socket>();
|
||||
this.swissTimingAdapterPersistence = swissTimingAdapterPersistence;
|
||||
this.swissTimingFactory = swissTimingFactory;
|
||||
lastMessageCountCollection = db.getCollection(CollectionNames.LAST_MESSAGE_COUNT.name());
|
||||
lastMessageCount = getLastMessageCount();
|
||||
clientListener = createClientListenerThread(portForClients);
|
||||
clientListener.start();
|
||||
storeAndForwardThread = new Thread(this, "StoreAndForward");
|
||||
storeAndForwardThread.start();
|
||||
synchronized (this) {
|
||||
while (!listeningForClients) {
|
||||
wait();
|
||||
}
|
||||
while (!receivingFromSailMaster) {
|
||||
wait();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates a storing message forwarder in listening mode. In this mode, this object won't actively try to open
|
||||
* TCP connections to a SailMaster system / bridge but instead listen for inbound TCP connections on port
|
||||
* <code>listenPort</code>.
|
||||
*
|
||||
* @param listenPort
|
||||
* listens on this port for messages coming in from a real SwissTiming SailMaster
|
||||
* @param portForClients
|
||||
* clients can connect to this port and will receive forwarded and sequence-numbered messages over those
|
||||
* sockets
|
||||
* @throws IOException
|
||||
*/
|
||||
public StoreAndForward(final int listenPort, final int portForClients, SwissTimingFactory swissTimingFactory,
|
||||
SwissTimingAdapterPersistence swissTimingAdapterPersistence, MongoDBService mongoDBService) throws InterruptedException, IOException {
|
||||
this.db = mongoDBService.getDB();
|
||||
this.listenPort = listenPort;
|
||||
this.transceiver = swissTimingFactory.createSailMasterTransceiver();
|
||||
this.portForClients = portForClients;
|
||||
this.streamsToForwardTo = new ArrayList<OutputStream>();
|
||||
this.socketsToForwardTo = new ArrayList<Socket>();
|
||||
this.swissTimingAdapterPersistence = swissTimingAdapterPersistence;
|
||||
this.swissTimingFactory = swissTimingFactory;
|
||||
this.sailMasterHostname = null;
|
||||
this.sailMasterPort = -1;
|
||||
lastMessageCountCollection = db.getCollection(CollectionNames.LAST_MESSAGE_COUNT.name());
|
||||
lastMessageCount = getLastMessageCount();
|
||||
serverSocketListeningForSailMasterBridge = new ServerSocket(listenPort);
|
||||
clientListener = createClientListenerThread(portForClients);
|
||||
clientListener.start();
|
||||
storeAndForwardThread = new Thread(this, "StoreAndForward");
|
||||
storeAndForwardThread.start();
|
||||
synchronized (this) {
|
||||
while (!listeningForClients) {
|
||||
wait();
|
||||
}
|
||||
while (!receivingFromSailMaster) {
|
||||
wait();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private long getLastMessageCount() {
|
||||
DBObject lastMessageCountRecord = lastMessageCountCollection.findOne();
|
||||
if (lastMessageCountRecord == null) {
|
||||
lastMessageCountCollection.insert(new BasicDBObject().append(FieldNames.LAST_MESSAGE_COUNT.name(), 0l));
|
||||
}
|
||||
return lastMessageCountRecord == null ? 0 : ((Number) lastMessageCountRecord.get(FieldNames.LAST_MESSAGE_COUNT.name())).longValue();
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns <code>true</code> if and only if this object is listening for incoming SailMaster TCP
|
||||
* connections instead of actively connecting / reconnecting to a SailMaster server by itself.
|
||||
*/
|
||||
private boolean isInSailMasterListeningMode() {
|
||||
return serverSocketListeningForSailMasterBridge != null;
|
||||
}
|
||||
|
||||
private Thread createClientListenerThread(final int portForClients) {
|
||||
return new Thread(new Runnable() {
|
||||
public void run() {
|
||||
ServerSocket ss;
|
||||
try {
|
||||
synchronized (StoreAndForward.this) {
|
||||
ss = new ServerSocket(portForClients);
|
||||
listeningForClients = true;
|
||||
logger.info("StoreAndForward listening for clients on port "+portForClients);
|
||||
StoreAndForward.this.notifyAll();
|
||||
}
|
||||
while (!stopped) {
|
||||
Socket s = ss.accept();
|
||||
logger.info("StoreAndForward received connector's connect request on port "+portForClients);
|
||||
if (!stopped) {
|
||||
synchronized (StoreAndForward.this) {
|
||||
socketsToForwardTo.add(s);
|
||||
streamsToForwardTo.add(s.getOutputStream());
|
||||
}
|
||||
} else {
|
||||
s.close();
|
||||
}
|
||||
}
|
||||
ss.close();
|
||||
logger.info("StoreAndForward client listener thread stopped.");
|
||||
} catch (IOException e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
}
|
||||
}, "StoreAndForwardClientListener");
|
||||
}
|
||||
|
||||
/**
|
||||
* Depending on the mode of operation, accepts an inbound connect from a SailMaster server / bridge or
|
||||
* actively initiates a TCP connection to the SailMaster address/port configured. When this method
|
||||
* returns, the {@link #socket} holds an open and connected socket.
|
||||
* @throws IOException
|
||||
*/
|
||||
private void establishConnection() throws IOException {
|
||||
if (isInSailMasterListeningMode()) {
|
||||
synchronized (this) {
|
||||
receivingFromSailMaster = true;
|
||||
logger.info("StoreAndForward waiting for inbound SailMaster connections on port "+listenPort);
|
||||
notifyAll();
|
||||
}
|
||||
socket = serverSocketListeningForSailMasterBridge.accept();
|
||||
logger.info("StoreAndForward received SailMaster connect on port "+listenPort);
|
||||
} else {
|
||||
synchronized (this) {
|
||||
receivingFromSailMaster = true;
|
||||
logger.info("StoreAndForward issuing SailMaster connection to "+sailMasterHostname+":"+sailMasterPort);
|
||||
notifyAll();
|
||||
}
|
||||
socket = new Socket(sailMasterHostname, sailMasterPort);
|
||||
logger.info("StoreAndForward connections to SailMaster "+sailMasterHostname+":"+sailMasterPort+" established");
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Stops execution after having received the next message
|
||||
*/
|
||||
public void stop() throws UnknownHostException, IOException, InterruptedException {
|
||||
logger.entering(getClass().getName(), "stop");
|
||||
stopped = true;
|
||||
new Socket("localhost", portForClients); // this is to stop the client listener thread
|
||||
logger.info("joining clientListener thread "+clientListener);
|
||||
clientListener.join();
|
||||
socket.close(); // will let a read terminate abnormally
|
||||
logger.info("joining storeAndForwardThread "+storeAndForwardThread);
|
||||
storeAndForwardThread.join();
|
||||
}
|
||||
|
||||
public static void main(String[] args) throws InterruptedException, IOException {
|
||||
String hostname = null;
|
||||
int i=0;
|
||||
if (args[i].equals("-h")) {
|
||||
hostname = args[++i];
|
||||
i++;
|
||||
}
|
||||
int sailMasterPort = Integer.valueOf(args[i++]);
|
||||
int clientPort = Integer.valueOf(args[i++]);
|
||||
|
||||
MongoDBService mongoDBService = MongoDBService.INSTANCE;
|
||||
mongoDBService.setConfiguration(MongoDBConfiguration.getDefaultConfiguration());
|
||||
SwissTimingAdapterPersistence swissTimingAdapterPersistence = SwissTimingAdapterPersistence.INSTANCE;
|
||||
if (hostname == null) {
|
||||
new StoreAndForward(sailMasterPort, clientPort, SwissTimingFactory.INSTANCE, swissTimingAdapterPersistence, mongoDBService);
|
||||
} else {
|
||||
new StoreAndForward(hostname, sailMasterPort, clientPort, SwissTimingFactory.INSTANCE, swissTimingAdapterPersistence, mongoDBService);
|
||||
}
|
||||
}
|
||||
|
||||
public void run() {
|
||||
logger.entering(getClass().getName(), "run");
|
||||
try {
|
||||
while (!stopped) {
|
||||
try {
|
||||
establishConnection();
|
||||
InputStream is = socket.getInputStream();
|
||||
Pair<String, Long> messageAndOptionalSequenceNumber = transceiver.receiveMessage(is);
|
||||
// ignore any sequence number contained in the message; we'll create our own
|
||||
DBObject emptyQuery = new BasicDBObject();
|
||||
DBObject incrementLastMessageCountQuery = new BasicDBObject().
|
||||
append("$inc", new BasicDBObject().append(FieldNames.LAST_MESSAGE_COUNT.name(), 1));
|
||||
while (!stopped && messageAndOptionalSequenceNumber != null) {
|
||||
logger.fine("Received message: "+messageAndOptionalSequenceNumber.getA());
|
||||
DBObject newCountRecord = lastMessageCountCollection.findAndModify(emptyQuery, incrementLastMessageCountQuery);
|
||||
lastMessageCount = ((newCountRecord == null) ? 0l :
|
||||
((Number) newCountRecord.get(FieldNames.LAST_MESSAGE_COUNT.name())).longValue());
|
||||
SailMasterMessage message = swissTimingFactory.createMessage(messageAndOptionalSequenceNumber.getA(), lastMessageCount);
|
||||
swissTimingAdapterPersistence.storeSailMasterMessage(message);
|
||||
synchronized (this) {
|
||||
for (OutputStream os : new ArrayList<OutputStream>(streamsToForwardTo)) {
|
||||
// write the sequence number of the message into the stream before actually writing the
|
||||
// SwissTiming message
|
||||
try {
|
||||
// TODO if forwarding to os doesn't work, e.g., because the socket was closed or the client died, remove os from streamsToForwardTo and the socket from socketsToForwardTo
|
||||
transceiver.sendMessage(message, os);
|
||||
} catch (Throwable e) {
|
||||
int i=streamsToForwardTo.indexOf(os);
|
||||
try {
|
||||
os.close();
|
||||
} catch (Throwable t) {
|
||||
logger.throwing(StoreAndForward.class.getName(), "run", t);
|
||||
}
|
||||
streamsToForwardTo.remove(os);
|
||||
Socket s = socketsToForwardTo.remove(i);
|
||||
try {
|
||||
s.close();
|
||||
} catch (Throwable t) {
|
||||
logger.throwing(StoreAndForward.class.getName(), "run", t);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if (!stopped) {
|
||||
messageAndOptionalSequenceNumber = transceiver.receiveMessage(is);
|
||||
}
|
||||
}
|
||||
for (OutputStream os : streamsToForwardTo) {
|
||||
os.close();
|
||||
}
|
||||
for (Socket socketToForwardTo : socketsToForwardTo) {
|
||||
socketToForwardTo.close();
|
||||
}
|
||||
} catch (Throwable e) {
|
||||
e.printStackTrace();
|
||||
if (!stopped) {
|
||||
logger.throwing(StoreAndForward.class.getName(), "Error during forwarding message. Continuing...", e);
|
||||
try {
|
||||
Thread.sleep(1000l); // wait a little bit before trying to re-establish a connection
|
||||
} catch (InterruptedException e1) {
|
||||
logger.throwing(StoreAndForward.class.getName(), "Can't find any sleep...", e1);
|
||||
}
|
||||
} else {
|
||||
logger.info("StoreAndForward socket was closed.");
|
||||
}
|
||||
}
|
||||
}
|
||||
logger.info("StoreAndForward is closing receiving server socket");
|
||||
if (serverSocketListeningForSailMasterBridge != null && !serverSocketListeningForSailMasterBridge.isClosed()) {
|
||||
serverSocketListeningForSailMasterBridge.close();
|
||||
}
|
||||
logger.info("Stopping StoreAndForward server.");
|
||||
} catch (IOException e) {
|
||||
logger.throwing(getClass().getName(), "run", e);
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
logger.exiting(getClass().getName(), "run");
|
||||
}
|
||||
}
|
||||
|
||||
+6
-6
@@ -1,6 +1,6 @@
|
||||
package com.sap.sailing.domain.swisstimingadapter.persistence.impl;
|
||||
|
||||
public enum CollectionNames {
|
||||
SWISSTIMING_CONFIGURATIONS, LAST_MESSAGE_COUNT,
|
||||
RACES_MASTERDATA, RACES_MESSAGES, COMMAND_MESSAGES
|
||||
}
|
||||
package com.sap.sailing.domain.swisstimingadapter.persistence.impl;
|
||||
|
||||
public enum CollectionNames {
|
||||
SWISSTIMING_CONFIGURATIONS, LAST_MESSAGE_COUNT,
|
||||
RACES_MASTERDATA, RACES_MESSAGES, COMMAND_MESSAGES
|
||||
}
|
||||
|
||||
+15
-15
@@ -1,15 +1,15 @@
|
||||
package com.sap.sailing.domain.swisstimingadapter.persistence.impl;
|
||||
|
||||
public enum FieldNames {
|
||||
// SwissTiming configuration parameters:
|
||||
ST_CONFIG_NAME, ST_CONFIG_HOSTNAME, ST_CONFIG_PORT, ST_CONFIG_CAN_SEND_REQUESTS,
|
||||
|
||||
// last message count field:
|
||||
LAST_MESSAGE_COUNT,
|
||||
|
||||
// raw messages:
|
||||
MESSAGE_SEQUENCE_NUMBER, MESSAGE_CONTENT, MESSAGE_COMMAND,
|
||||
|
||||
// race specific message and masterdata
|
||||
RACE_ID, RACE_DESCRIPTION, RACE_STARTTIME,
|
||||
}
|
||||
package com.sap.sailing.domain.swisstimingadapter.persistence.impl;
|
||||
|
||||
public enum FieldNames {
|
||||
// SwissTiming configuration parameters:
|
||||
ST_CONFIG_NAME, ST_CONFIG_HOSTNAME, ST_CONFIG_PORT, ST_CONFIG_CAN_SEND_REQUESTS,
|
||||
|
||||
// last message count field:
|
||||
LAST_MESSAGE_COUNT,
|
||||
|
||||
// raw messages:
|
||||
MESSAGE_SEQUENCE_NUMBER, MESSAGE_CONTENT, MESSAGE_COMMAND,
|
||||
|
||||
// race specific message and masterdata
|
||||
RACE_ID, RACE_DESCRIPTION, RACE_STARTTIME,
|
||||
}
|
||||
|
||||
+2
-1
@@ -117,6 +117,7 @@ public class SwissTimingAdapterPersistenceImpl implements SwissTimingAdapterPers
|
||||
@Override
|
||||
public List<SailMasterMessage> loadRaceMessages(String raceID) {
|
||||
DBCollection racesMessagesCollection = database.getCollection(CollectionNames.RACES_MESSAGES.name());
|
||||
racesMessagesCollection.ensureIndex(new BasicDBObject(FieldNames.MESSAGE_SEQUENCE_NUMBER.name(), null)); // no sort without index
|
||||
BasicDBObject query = new BasicDBObject();
|
||||
query.append(FieldNames.RACE_ID.name(), raceID);
|
||||
DBCursor results = racesMessagesCollection.find(query).sort(
|
||||
@@ -124,7 +125,7 @@ public class SwissTimingAdapterPersistenceImpl implements SwissTimingAdapterPers
|
||||
List<SailMasterMessage> result = new ArrayList<SailMasterMessage>();
|
||||
for (DBObject o : results) {
|
||||
SailMasterMessage msg = swissTimingFactory.createMessage((String) o.get(FieldNames.MESSAGE_CONTENT.name()),
|
||||
(Long) o.get(FieldNames.MESSAGE_SEQUENCE_NUMBER.name()));
|
||||
((Number) o.get(FieldNames.MESSAGE_SEQUENCE_NUMBER.name())).longValue());
|
||||
result.add(msg);
|
||||
}
|
||||
return result;
|
||||
|
||||
Reference in New Issue
Block a user