first successful test with numbering forwarded SailMaster messages

This commit is contained in:
Axel Uhl committed 2011-11-09 12:57:38 +01:00
1 parent 46092d9321
commit cc6d9f330c
12 files changed
+243 -49

No files matched your search

@@ -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"
@@ -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<Pair<String, Integer>> hostnamesAndPorts;
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 final DBCollection lastMessageCountCollection;
public StoreAndForward(int listenPort, List<Pair<String, Integer>> 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<OutputStream>();
this.socketsToForwardTo = new ArrayList<Socket>();
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<Pair<String, Integer>> hostnamesAndPorts = new ArrayList<Pair<String,Integer>>();
for (int i=1; i<args.length-1; i+=2) {
Pair<String, Integer> hostnameAndPort = new Pair<String, Integer>(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<OutputStream> streamsToForwardTo = new ArrayList<OutputStream>();
List<Socket> socketsToForwardTo = new ArrayList<Socket>();
for (Pair<String, Integer> 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<String, Integer> 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);
}
@@ -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/,
.
@@ -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<Race> racesReceived = new ArrayList<Race>();
final boolean[] receivedSomething = new boolean[1];
connector.addSailMasterListener(new SailMasterAdapter() {
@Override
public void receivedAvailableRaces(Iterable<Race> 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);
}
}
@@ -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);
}
@@ -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<Fix> fixes) {
}
@Override
public void receivedTimingData(String raceID, String boatID,
List<Triple<Integer, Integer, Long>> markIndicesRanksAndTimesSinceStartInMilliseconds) {
}
@Override
public void receivedClockAtMark(String raceID,
List<Triple<Integer, TimePoint, String>> markIndicesTimePointsAndBoatIDs) {
}
@Override
public void receivedStartList(String raceID, StartList startList) {
}
@Override
public void receivedCourseConfiguration(String raceID, Course course) {
}
@Override
public void receivedAvailableRaces(Iterable<Race> races) {
}
}
@@ -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;
}
@@ -69,6 +69,7 @@ public class SailMasterConnectorImpl extends SailMasterTransceiverImpl implement
private final Set<SailMasterListener> 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<MessageType, BlockingQueue<SailMasterMessage>> 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<String, Integer> 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();
}
}
}
@@ -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);
}
@@ -18,7 +18,7 @@ public class SwissTimingRaceTrackerImpl implements SwissTimingRaceTracker {
private final SailMasterConnector connector;
// private final Set<RaceDefinition> 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<RaceDefinition>();
}
@@ -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();
@@ -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<String, String, Integer> key = new Triple<String, String, Integer>(raceID, hostname, port);
RaceTracker tracker = raceTrackersByID.get(key);
if (tracker == null) {