diff --git a/java/com.sap.sailing.domain.swisstimingadapter.persistence/META-INF/MANIFEST.MF b/java/com.sap.sailing.domain.swisstimingadapter.persistence/META-INF/MANIFEST.MF index 07f856231e0..3c9f2d4de01 100755 --- a/java/com.sap.sailing.domain.swisstimingadapter.persistence/META-INF/MANIFEST.MF +++ b/java/com.sap.sailing.domain.swisstimingadapter.persistence/META-INF/MANIFEST.MF @@ -9,4 +9,5 @@ Require-Bundle: com.sap.sailing.domain.swisstimingadapter, com.mongodb.driver;bundle-version="2.6.2", com.sap.sailing.mongodb, com.sap.sailing.domain -Export-Package: com.sap.sailing.domain.swisstimingadapter.persistence +Export-Package: com.sap.sailing.domain.swisstimingadapter.persistence, + com.sap.sailing.domain.swisstimingadapter.persistence.impl;x-friends:="com.sap.sailing.domain.swisstimingadapter.test" diff --git a/java/com.sap.sailing.domain.swisstimingadapter.persistence/src/com/sap/sailing/domain/swisstimingadapter/persistence/StoreAndForward.java b/java/com.sap.sailing.domain.swisstimingadapter.persistence/src/com/sap/sailing/domain/swisstimingadapter/persistence/StoreAndForward.java index 052cad1ea55..29f266f605d 100755 --- a/java/com.sap.sailing.domain.swisstimingadapter.persistence/src/com/sap/sailing/domain/swisstimingadapter/persistence/StoreAndForward.java +++ b/java/com.sap.sailing.domain.swisstimingadapter.persistence/src/com/sap/sailing/domain/swisstimingadapter/persistence/StoreAndForward.java @@ -5,6 +5,7 @@ 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; @@ -36,50 +37,89 @@ public class StoreAndForward implements Runnable { private final DB db; private final int listenPort; private final SailMasterTransceiver transceiver; - private final List> hostnamesAndPorts; + private final int portForClients; private long lastMessageCount; + private boolean stopped; + private final Thread clientListener; + private final List socketsToForwardTo; + private final List streamsToForwardTo; + private boolean listeningForClients; private final DBCollection lastMessageCountCollection; - public StoreAndForward(int listenPort, List> hostnamesAndPorts) { + /** + * @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 + */ + public StoreAndForward(int listenPort, final int portForClients) throws InterruptedException { db = Activator.getDefaultInstance().getDB(); this.listenPort = listenPort; this.transceiver = SwissTimingFactory.INSTANCE.createSailMasterTransceiver(); - this.hostnamesAndPorts = hostnamesAndPorts; + this.portForClients = portForClients; + this.streamsToForwardTo = new ArrayList(); + this.socketsToForwardTo = new ArrayList(); lastMessageCountCollection = db.getCollection(CollectionNames.LAST_MESSAGE_COUNT.name()); DBObject lastMessageCountRecord = lastMessageCountCollection.findOne(); - lastMessageCount = (Long) lastMessageCountRecord.get(FieldNames.LAST_MESSAGE_COUNT.name()); + lastMessageCount = lastMessageCountRecord == null ? 0 : (Long) lastMessageCountRecord.get(FieldNames.LAST_MESSAGE_COUNT.name()); + clientListener = new Thread(new Runnable() { + public void run() { + ServerSocket ss; + try { + synchronized (StoreAndForward.this) { + ss = new ServerSocket(portForClients); + listeningForClients = true; + StoreAndForward.this.notifyAll(); + } + while (!stopped) { + Socket s = ss.accept(); + if (!stopped) { + synchronized (StoreAndForward.this) { + socketsToForwardTo.add(s); + streamsToForwardTo.add(s.getOutputStream()); + } + } else { + s.close(); + } + } + logger.info("StoreAndForward client listener thread stopped."); + } catch (IOException e) { + throw new RuntimeException(e); + } + } + }, "StoreAndForwardClientListener"); + clientListener.start(); + synchronized (this) { + while (!listeningForClients) { + wait(); + } + } + } + + /** + * Stops execution after having received the next message + */ + public void stop() throws UnknownHostException, IOException, InterruptedException { + stopped = true; + new Socket("localhost", portForClients); // this is to stop the client listener thread + clientListener.join(); } - public static void main(String[] args) { + public static void main(String[] args) throws InterruptedException { int listenPort = Integer.valueOf(args[0]); - List> hostnamesAndPorts = new ArrayList>(); - for (int i=1; i hostnameAndPort = new Pair(args[i], Integer.valueOf(args[i+1])); - hostnamesAndPorts.add(hostnameAndPort); - } - StoreAndForward storeAndForward = new StoreAndForward(listenPort, hostnamesAndPorts); + int clientPort = Integer.valueOf(args[1]); + StoreAndForward storeAndForward = new StoreAndForward(listenPort, clientPort); storeAndForward.run(); } public void run() { try { ServerSocket ss = new ServerSocket(listenPort); - while (true) { + while (!stopped) { Socket socket = ss.accept(); try { - List streamsToForwardTo = new ArrayList(); - List socketsToForwardTo = new ArrayList(); - for (Pair hostnameAndPort : hostnamesAndPorts) { - try { - Socket socketToForwardTo = new Socket(hostnameAndPort.getA(), hostnameAndPort.getB()); - socketsToForwardTo.add(socketToForwardTo); - streamsToForwardTo.add(socketToForwardTo.getOutputStream()); - } catch (Exception e) { - logger.throwing(StoreAndForward.class.getName(), "While trying to open a socket to forward to "+ - hostnameAndPort, e); - } - } InputStream is = socket.getInputStream(); Pair messageAndOptionalSequenceNumber = transceiver.receiveMessage(is); // ignore any sequence number contained in the message; we'll create our own @@ -89,12 +129,17 @@ public class StoreAndForward implements Runnable { while (messageAndOptionalSequenceNumber != null) { DBObject newCountRecord = lastMessageCountCollection.findAndModify(emptyQuery, incrementLastMessageCountQuery); lastMessageCount = (Long) newCountRecord.get(FieldNames.LAST_MESSAGE_COUNT.name()); - for (OutputStream os : streamsToForwardTo) { - // write the sequence number of the message into the stream before actually writing the SwissTiming message - os.write((""+lastMessageCount).getBytes()); - transceiver.sendMessage(messageAndOptionalSequenceNumber.getA(), os); + synchronized (this) { + for (OutputStream os : streamsToForwardTo) { + // write the sequence number of the message into the stream before actually writing the + // SwissTiming message + os.write(("" + lastMessageCount).getBytes()); + transceiver.sendMessage(messageAndOptionalSequenceNumber.getA(), os); + } + } + if (!stopped) { + messageAndOptionalSequenceNumber = transceiver.receiveMessage(is); } - messageAndOptionalSequenceNumber = transceiver.receiveMessage(is); } for (OutputStream os : streamsToForwardTo) { os.close(); @@ -103,9 +148,10 @@ public class StoreAndForward implements Runnable { socketToForwardTo.close(); } } catch (Exception e) { - + logger.throwing(StoreAndForward.class.getName(), "Error during forwarding message. Continuing...", e); } } + logger.info("Stopping StoreAndForward server."); } catch (IOException e) { throw new RuntimeException(e); } diff --git a/java/com.sap.sailing.domain.swisstimingadapter.test/META-INF/MANIFEST.MF b/java/com.sap.sailing.domain.swisstimingadapter.test/META-INF/MANIFEST.MF index 66929b12d6e..1008467c686 100755 --- a/java/com.sap.sailing.domain.swisstimingadapter.test/META-INF/MANIFEST.MF +++ b/java/com.sap.sailing.domain.swisstimingadapter.test/META-INF/MANIFEST.MF @@ -8,6 +8,9 @@ Bundle-RequiredExecutionEnvironment: JavaSE-1.6 Require-Bundle: com.sap.sailing.domain.swisstimingadapter, com.sap.sailing.domain, org.junit4, - com.sap.sailing.udpconnector + com.sap.sailing.udpconnector, + com.sap.sailing.domain.swisstimingadapter.persistence, + com.mongodb.driver;bundle-version="2.6.2", + com.sap.sailing.mongodb Bundle-ClassPath: resources/, . diff --git a/java/com.sap.sailing.domain.swisstimingadapter.test/src/com/sap/sailing/domain/swisstimingadapter/test/StoreAndForwardTest.java b/java/com.sap.sailing.domain.swisstimingadapter.test/src/com/sap/sailing/domain/swisstimingadapter/test/StoreAndForwardTest.java new file mode 100755 index 00000000000..358c3e04c2c --- /dev/null +++ b/java/com.sap.sailing.domain.swisstimingadapter.test/src/com/sap/sailing/domain/swisstimingadapter/test/StoreAndForwardTest.java @@ -0,0 +1,93 @@ +package com.sap.sailing.domain.swisstimingadapter.test; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; + +import java.io.IOException; +import java.io.OutputStream; +import java.net.Socket; +import java.net.UnknownHostException; +import java.util.ArrayList; +import java.util.List; + +import org.junit.After; +import org.junit.Before; +import org.junit.Test; + +import com.mongodb.BasicDBObject; +import com.mongodb.DB; +import com.mongodb.DBCollection; +import com.sap.sailing.domain.swisstimingadapter.MessageType; +import com.sap.sailing.domain.swisstimingadapter.Race; +import com.sap.sailing.domain.swisstimingadapter.SailMasterAdapter; +import com.sap.sailing.domain.swisstimingadapter.SailMasterConnector; +import com.sap.sailing.domain.swisstimingadapter.SailMasterTransceiver; +import com.sap.sailing.domain.swisstimingadapter.SwissTimingFactory; +import com.sap.sailing.domain.swisstimingadapter.persistence.StoreAndForward; +import com.sap.sailing.domain.swisstimingadapter.persistence.impl.CollectionNames; +import com.sap.sailing.domain.swisstimingadapter.persistence.impl.FieldNames; +import com.sap.sailing.mongodb.Activator; + +public class StoreAndForwardTest { + private static final int RECEIVE_PORT = 6543; + private static final int CLIENT_PORT = 6544; + + private DB db; + private StoreAndForward storeAndForward; + private Thread storeAndForwardThread; + private Socket sendingSocket; + private OutputStream sendingStream; + private SailMasterTransceiver transceiver; + private SailMasterConnector connector; + + @Before + public void setUp() throws UnknownHostException, IOException, InterruptedException { + db = Activator.getDefaultInstance().getDB(); + storeAndForward = new StoreAndForward(RECEIVE_PORT, CLIENT_PORT); + storeAndForwardThread = new Thread(storeAndForward, "StoreAndForward"); + storeAndForwardThread.start(); + sendingSocket = new Socket("localhost", RECEIVE_PORT); + sendingStream = sendingSocket.getOutputStream(); + transceiver = SwissTimingFactory.INSTANCE.createSailMasterTransceiver(); + connector = SwissTimingFactory.INSTANCE.createSailMasterConnector("localhost", CLIENT_PORT); + DBCollection lastMessageCountCollection = db.getCollection(CollectionNames.LAST_MESSAGE_COUNT.name()); + lastMessageCountCollection.update(new BasicDBObject(), new BasicDBObject().append(FieldNames.LAST_MESSAGE_COUNT.name(), 0l), + /* upsert */ true, /* multi */ false); + } + + @After + public void tearDown() throws InterruptedException, IOException { + storeAndForward.stop(); + transceiver.sendMessage(MessageType._STOPSERVER.name(), sendingStream); // forwards to connector and hence stops it + storeAndForwardThread.join(); + } + + @Test + public void testSimpleRACMessage() throws IOException, InterruptedException { + final List racesReceived = new ArrayList(); + final boolean[] receivedSomething = new boolean[1]; + connector.addSailMasterListener(new SailMasterAdapter() { + @Override + public void receivedAvailableRaces(Iterable races) { + for (Race race : races) { + racesReceived.add(race); + } + synchronized (StoreAndForwardTest.this) { + receivedSomething[0] = true; + StoreAndForwardTest.this.notifyAll(); + } + } + }); + transceiver.sendMessage("RAC|2|4711;A wonderful test race|4712;Not such a wonderful race", sendingStream); + synchronized (this) { + if (!receivedSomething[0]) { + wait(2000l); // wait for two seconds to receive the message + } + } + assertTrue(receivedSomething[0]); + assertEquals(2, racesReceived.size()); + DBCollection lastMessageCountCollection = db.getCollection(CollectionNames.LAST_MESSAGE_COUNT.name()); + Long lastMessageCount = (Long) lastMessageCountCollection.findOne().get(FieldNames.LAST_MESSAGE_COUNT.name()); + assertEquals((Long) 1l, lastMessageCount); + } +} diff --git a/java/com.sap.sailing.domain.swisstimingadapter.test/src/com/sap/sailing/domain/swisstimingadapter/test/SwissTimingSailMasterLiveTest.java b/java/com.sap.sailing.domain.swisstimingadapter.test/src/com/sap/sailing/domain/swisstimingadapter/test/SwissTimingSailMasterLiveTest.java index 9c32ad29627..11f9fc804db 100755 --- a/java/com.sap.sailing.domain.swisstimingadapter.test/src/com/sap/sailing/domain/swisstimingadapter/test/SwissTimingSailMasterLiveTest.java +++ b/java/com.sap.sailing.domain.swisstimingadapter.test/src/com/sap/sailing/domain/swisstimingadapter/test/SwissTimingSailMasterLiveTest.java @@ -41,7 +41,7 @@ public class SwissTimingSailMasterLiveTest implements SailMasterListener { private SailMasterConnector connector; @Before - public void connect() { + public void connect() throws InterruptedException { connector = SwissTimingFactory.INSTANCE.createSailMasterConnector("gps.sportresult.com", 40300); } diff --git a/java/com.sap.sailing.domain.swisstimingadapter/src/com/sap/sailing/domain/swisstimingadapter/SailMasterAdapter.java b/java/com.sap.sailing.domain.swisstimingadapter/src/com/sap/sailing/domain/swisstimingadapter/SailMasterAdapter.java new file mode 100755 index 00000000000..e7a0d8e68e5 --- /dev/null +++ b/java/com.sap.sailing.domain.swisstimingadapter/src/com/sap/sailing/domain/swisstimingadapter/SailMasterAdapter.java @@ -0,0 +1,40 @@ +package com.sap.sailing.domain.swisstimingadapter; + +import java.util.Collection; +import java.util.List; + +import com.sap.sailing.domain.base.Distance; +import com.sap.sailing.domain.base.TimePoint; +import com.sap.sailing.util.Util.Triple; + +public abstract class SailMasterAdapter implements SailMasterListener { + + @Override + public void receivedRacePositionData(String raceID, RaceStatus status, TimePoint timePoint, TimePoint startTime, + Long millisecondsSinceRaceStart, Integer nextMarkIndexForLeader, Distance distanceToNextMarkForLeader, + Collection fixes) { + } + + @Override + public void receivedTimingData(String raceID, String boatID, + List> markIndicesRanksAndTimesSinceStartInMilliseconds) { + } + + @Override + public void receivedClockAtMark(String raceID, + List> markIndicesTimePointsAndBoatIDs) { + } + + @Override + public void receivedStartList(String raceID, StartList startList) { + } + + @Override + public void receivedCourseConfiguration(String raceID, Course course) { + } + + @Override + public void receivedAvailableRaces(Iterable races) { + } + +} diff --git a/java/com.sap.sailing.domain.swisstimingadapter/src/com/sap/sailing/domain/swisstimingadapter/SwissTimingFactory.java b/java/com.sap.sailing.domain.swisstimingadapter/src/com/sap/sailing/domain/swisstimingadapter/SwissTimingFactory.java index e8f3bd62378..7a5ac97376e 100755 --- a/java/com.sap.sailing.domain.swisstimingadapter/src/com/sap/sailing/domain/swisstimingadapter/SwissTimingFactory.java +++ b/java/com.sap.sailing.domain.swisstimingadapter/src/com/sap/sailing/domain/swisstimingadapter/SwissTimingFactory.java @@ -8,11 +8,11 @@ public interface SwissTimingFactory { SwissTimingMessageParser createMessageParser(); - SailMasterConnector createSailMasterConnector(String hostname, int port); + SailMasterConnector createSailMasterConnector(String hostname, int port) throws InterruptedException; SailMasterTransceiver createSailMasterTransceiver(); SwissTimingConfiguration createSwissTimingConfiguration(String name, String hostname, int port); - SwissTimingRaceTracker createRaceTracker(String raceID, String hostname, int port, WindStore windStore); + SwissTimingRaceTracker createRaceTracker(String raceID, String hostname, int port, WindStore windStore) throws InterruptedException; } diff --git a/java/com.sap.sailing.domain.swisstimingadapter/src/com/sap/sailing/domain/swisstimingadapter/impl/SailMasterConnectorImpl.java b/java/com.sap.sailing.domain.swisstimingadapter/src/com/sap/sailing/domain/swisstimingadapter/impl/SailMasterConnectorImpl.java index df06b86a055..f5bf034474e 100755 --- a/java/com.sap.sailing.domain.swisstimingadapter/src/com/sap/sailing/domain/swisstimingadapter/impl/SailMasterConnectorImpl.java +++ b/java/com.sap.sailing.domain.swisstimingadapter/src/com/sap/sailing/domain/swisstimingadapter/impl/SailMasterConnectorImpl.java @@ -69,6 +69,7 @@ public class SailMasterConnectorImpl extends SailMasterTransceiverImpl implement private final Set listeners; private final Thread receiverThread; private boolean stopped; + private boolean connected; /** * Currently the SwissTiming SailMaster protocol only transmits time zone information when sending @@ -85,7 +86,7 @@ public class SailMasterConnectorImpl extends SailMasterTransceiverImpl implement private final Map> unprocessedMessagesByType; - public SailMasterConnectorImpl(String host, int port) { + public SailMasterConnectorImpl(String host, int port) throws InterruptedException { super(); dateFormat = new SimpleDateFormat("yyyy-MM-dd'T'hh:mm:ssZ"); this.host = host; @@ -96,6 +97,11 @@ public class SailMasterConnectorImpl extends SailMasterTransceiverImpl implement lastTimeZoneSuffix = (offset<0?"-":"+") + new DecimalFormat("00").format(offset)+"00"; receiverThread = new Thread(this, "SwissTiming SailMaster Receiver"); receiverThread.start(); + synchronized (this) { + while (!connected) { + wait(); + } + } } public void run() { @@ -106,15 +112,16 @@ public class SailMasterConnectorImpl extends SailMasterTransceiverImpl implement Pair receivedMessageAndOptionalSequenceNumber = receiveMessage(socket.getInputStream()); SailMasterMessage message = new SailMasterMessageImpl(receivedMessageAndOptionalSequenceNumber.getA(), receivedMessageAndOptionalSequenceNumber.getB()); - if (message.isResponse()) { - // this is a response for an explicit request - rendevouz(message); - if (message.getType() == MessageType._STOPSERVER) { - stop(); + if (message.getType() == MessageType._STOPSERVER) { + stop(); + } else { + if (message.isResponse()) { + // this is a response for an explicit request + rendevouz(message); + } else if (message.isEvent()) { + // a spontaneous event + notifyListeners(message); } - } else if (message.isEvent()) { - // a spontaneous event - notifyListeners(message); } } catch (SocketException se) { // This occurs if the socket was closed which may mean the connector was stopped. Check in while @@ -311,6 +318,10 @@ public class SailMasterConnectorImpl extends SailMasterTransceiverImpl implement private synchronized void ensureSocketIsOpen() throws UnknownHostException, IOException { if (socket == null) { socket = new Socket(host, port); + synchronized (this) { + connected = true; + notifyAll(); + } } } diff --git a/java/com.sap.sailing.domain.swisstimingadapter/src/com/sap/sailing/domain/swisstimingadapter/impl/SwissTimingFactoryImpl.java b/java/com.sap.sailing.domain.swisstimingadapter/src/com/sap/sailing/domain/swisstimingadapter/impl/SwissTimingFactoryImpl.java index 1a7e65f6a08..15e202f9192 100755 --- a/java/com.sap.sailing.domain.swisstimingadapter/src/com/sap/sailing/domain/swisstimingadapter/impl/SwissTimingFactoryImpl.java +++ b/java/com.sap.sailing.domain.swisstimingadapter/src/com/sap/sailing/domain/swisstimingadapter/impl/SwissTimingFactoryImpl.java @@ -16,7 +16,7 @@ public class SwissTimingFactoryImpl implements SwissTimingFactory { } @Override - public SailMasterConnector createSailMasterConnector(String host, int port) { + public SailMasterConnector createSailMasterConnector(String host, int port) throws InterruptedException { return new SailMasterConnectorImpl(host, port); } @@ -26,7 +26,7 @@ public class SwissTimingFactoryImpl implements SwissTimingFactory { } @Override - public SwissTimingRaceTracker createRaceTracker(String raceID, String hostname, int port, WindStore windStore) { + public SwissTimingRaceTracker createRaceTracker(String raceID, String hostname, int port, WindStore windStore) throws InterruptedException { return new SwissTimingRaceTrackerImpl(raceID, hostname, port, this); } diff --git a/java/com.sap.sailing.domain.swisstimingadapter/src/com/sap/sailing/domain/swisstimingadapter/impl/SwissTimingRaceTrackerImpl.java b/java/com.sap.sailing.domain.swisstimingadapter/src/com/sap/sailing/domain/swisstimingadapter/impl/SwissTimingRaceTrackerImpl.java index 9f3d375920e..0c5b50d2bd8 100755 --- a/java/com.sap.sailing.domain.swisstimingadapter/src/com/sap/sailing/domain/swisstimingadapter/impl/SwissTimingRaceTrackerImpl.java +++ b/java/com.sap.sailing.domain.swisstimingadapter/src/com/sap/sailing/domain/swisstimingadapter/impl/SwissTimingRaceTrackerImpl.java @@ -18,7 +18,7 @@ public class SwissTimingRaceTrackerImpl implements SwissTimingRaceTracker { private final SailMasterConnector connector; // private final Set races; - protected SwissTimingRaceTrackerImpl(String raceID, String hostname, int port, SwissTimingFactory factory) { + protected SwissTimingRaceTrackerImpl(String raceID, String hostname, int port, SwissTimingFactory factory) throws InterruptedException { connector = factory.createSailMasterConnector(hostname, port); // races = new HashSet(); } diff --git a/java/com.sap.sailing.server/src/com/sap/sailing/server/RacingEventService.java b/java/com.sap.sailing.server/src/com/sap/sailing/server/RacingEventService.java index c9ffde1a1cf..45e130b3e34 100755 --- a/java/com.sap.sailing.server/src/com/sap/sailing/server/RacingEventService.java +++ b/java/com.sap.sailing.server/src/com/sap/sailing/server/RacingEventService.java @@ -176,7 +176,7 @@ public interface RacingEventService { */ void updateStoredLeaderboard(Leaderboard leaderboard); - RaceHandle addSwissTimingRace(String raceID, String hostname, int port, WindStore windStore, long timeoutInMilliseconds); + RaceHandle addSwissTimingRace(String raceID, String hostname, int port, WindStore windStore, long timeoutInMilliseconds) throws InterruptedException; SwissTimingFactory getSwissTimingFactory(); diff --git a/java/com.sap.sailing.server/src/com/sap/sailing/server/RacingEventServiceImpl.java b/java/com.sap.sailing.server/src/com/sap/sailing/server/RacingEventServiceImpl.java index 55a40759690..2e7a28cf142 100755 --- a/java/com.sap.sailing.server/src/com/sap/sailing/server/RacingEventServiceImpl.java +++ b/java/com.sap.sailing.server/src/com/sap/sailing/server/RacingEventServiceImpl.java @@ -220,7 +220,7 @@ public class RacingEventServiceImpl implements RacingEventService { } @Override - public synchronized RaceHandle addSwissTimingRace(String raceID, String hostname, int port, WindStore windStore, long timeoutInMilliseconds) { + public synchronized RaceHandle addSwissTimingRace(String raceID, String hostname, int port, WindStore windStore, long timeoutInMilliseconds) throws InterruptedException { Triple key = new Triple(raceID, hostname, port); RaceTracker tracker = raceTrackersByID.get(key); if (tracker == null) {