mirror of
https://github.com/eclipse-sailing-analytics/sailing-analytics.git
synced 2026-10-03 19:03:52 +00:00
first seemingly working test case showing the general mixing of messages loaded / received
This commit is contained in:
1 parent
0231120601
commit
a5ee6d0a2f
4 files changed
+84
-5
No files matched your search
+2
@@ -98,6 +98,7 @@ public class StoreAndForward implements Runnable {
|
||||
s.close();
|
||||
}
|
||||
}
|
||||
ss.close();
|
||||
logger.info("StoreAndForward client listener thread stopped.");
|
||||
} catch (IOException e) {
|
||||
throw new RuntimeException(e);
|
||||
@@ -181,6 +182,7 @@ public class StoreAndForward implements Runnable {
|
||||
}
|
||||
}
|
||||
}
|
||||
ss.close();
|
||||
logger.info("Stopping StoreAndForward server.");
|
||||
} catch (IOException e) {
|
||||
throw new RuntimeException(e);
|
||||
|
||||
+79
-3
@@ -7,6 +7,7 @@ import java.io.IOException;
|
||||
import java.io.OutputStream;
|
||||
import java.net.Socket;
|
||||
import java.net.UnknownHostException;
|
||||
import java.text.ParseException;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
|
||||
@@ -17,7 +18,10 @@ import org.junit.Test;
|
||||
import com.mongodb.BasicDBObject;
|
||||
import com.mongodb.DB;
|
||||
import com.mongodb.DBCollection;
|
||||
import com.sap.sailing.domain.swisstimingadapter.Course;
|
||||
import com.sap.sailing.domain.swisstimingadapter.Mark;
|
||||
import com.sap.sailing.domain.swisstimingadapter.Race;
|
||||
import com.sap.sailing.domain.swisstimingadapter.RaceSpecificMessageLoader;
|
||||
import com.sap.sailing.domain.swisstimingadapter.SailMasterAdapter;
|
||||
import com.sap.sailing.domain.swisstimingadapter.SailMasterConnector;
|
||||
import com.sap.sailing.domain.swisstimingadapter.SailMasterMessage;
|
||||
@@ -30,7 +34,7 @@ import com.sap.sailing.domain.swisstimingadapter.persistence.impl.CollectionName
|
||||
import com.sap.sailing.domain.swisstimingadapter.persistence.impl.FieldNames;
|
||||
import com.sap.sailing.mongodb.Activator;
|
||||
|
||||
public class StoreAndForwardTest {
|
||||
public class StoreAndForwardTest implements RaceSpecificMessageLoader {
|
||||
private static final int RECEIVE_PORT = 6543;
|
||||
private static final int CLIENT_PORT = 6544;
|
||||
|
||||
@@ -41,6 +45,9 @@ public class StoreAndForwardTest {
|
||||
private SailMasterTransceiver transceiver;
|
||||
private SailMasterConnector connector;
|
||||
private DomainObjectFactory domainObjectFactory;
|
||||
private ArrayList<SailMasterMessage> messagesToLoad;
|
||||
private SwissTimingFactory swissTimingFactory;
|
||||
private boolean loadMessagesCalled;
|
||||
|
||||
@Before
|
||||
public void setUp() throws UnknownHostException, IOException, InterruptedException {
|
||||
@@ -48,14 +55,16 @@ public class StoreAndForwardTest {
|
||||
storeAndForward = new StoreAndForward(RECEIVE_PORT, CLIENT_PORT, MongoObjectFactory.INSTANCE, SwissTimingFactory.INSTANCE);
|
||||
sendingSocket = new Socket("localhost", RECEIVE_PORT);
|
||||
sendingStream = sendingSocket.getOutputStream();
|
||||
transceiver = SwissTimingFactory.INSTANCE.createSailMasterTransceiver();
|
||||
connector = SwissTimingFactory.INSTANCE.createSailMasterConnector("localhost", CLIENT_PORT, domainObjectFactory);
|
||||
swissTimingFactory = SwissTimingFactory.INSTANCE;
|
||||
transceiver = swissTimingFactory.createSailMasterTransceiver();
|
||||
connector = swissTimingFactory.createSailMasterConnector("localhost", CLIENT_PORT, this);
|
||||
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);
|
||||
DBCollection rawMessages = db.getCollection(CollectionNames.RAW_MESSAGES.name());
|
||||
rawMessages.drop();
|
||||
domainObjectFactory = DomainObjectFactory.INSTANCE;
|
||||
messagesToLoad = new ArrayList<SailMasterMessage>();
|
||||
}
|
||||
|
||||
@After
|
||||
@@ -96,4 +105,71 @@ public class StoreAndForwardTest {
|
||||
assertEquals(rawMessage, rawMessages.get(0).getMessage());
|
||||
assertEquals(0l, (long) rawMessages.get(0).getSequenceNumber());
|
||||
}
|
||||
|
||||
/**
|
||||
* A {@link SailMasterConnector} can load messages from a {@link RaceSpecificMessageLoader} object when the
|
||||
* {@link SailMasterConnector#trackRace(String)} method is called. While it does so, messages received for the same
|
||||
* tracked race are buffered. When the loading of messages has finished, both, the loaded and buffered received
|
||||
* messages are parsed and notified to the listeners.<p>
|
||||
*
|
||||
* This test asserts that the general buffering mechanism works and that both, stored and sent messages are
|
||||
* notified properly.
|
||||
*/
|
||||
@Test
|
||||
public void testBuffering() throws UnknownHostException, IOException, InterruptedException, ParseException {
|
||||
final List<Course> coursesReceived = new ArrayList<Course>();
|
||||
final boolean[] receivedSomething = new boolean[1];
|
||||
connector.addSailMasterListener(new SailMasterAdapter() {
|
||||
@Override
|
||||
public void receivedCourseConfiguration(String raceID, Course course) {
|
||||
coursesReceived.add(course);
|
||||
synchronized (StoreAndForwardTest.this) {
|
||||
receivedSomething[0] = true;
|
||||
StoreAndForwardTest.this.notifyAll();
|
||||
}
|
||||
}
|
||||
});
|
||||
final int COUNT = 10;
|
||||
final int STORED = 5;
|
||||
assert COUNT <= 10 && STORED < COUNT;
|
||||
String[] rawMessage = new String[COUNT];
|
||||
for (int i=0; i<COUNT; i++) {
|
||||
rawMessage[i] = "CCG|4711|2|1;Lee Gate;LG1;LG2|"+i+";Windward;WW1";
|
||||
if (i<STORED) {
|
||||
messagesToLoad.add(swissTimingFactory.createMessage(rawMessage[i], (long) i));
|
||||
}
|
||||
}
|
||||
connector.trackRace("4711"); // this should transitively invoke loadMessages
|
||||
assertTrue(loadMessagesCalled);
|
||||
// for now, don't create an overlap; an overlap should be tested separately
|
||||
for (int i=STORED; i<COUNT; i++) {
|
||||
transceiver.sendMessage(rawMessage[i], sendingStream);
|
||||
}
|
||||
synchronized (this) {
|
||||
int attempts = 0;
|
||||
while (coursesReceived.size() < COUNT && attempts++ < 2 * COUNT) {
|
||||
wait(2000l); // wait for two seconds to receive the message
|
||||
}
|
||||
}
|
||||
assertTrue(receivedSomething[0]);
|
||||
assertEquals(COUNT, coursesReceived.size());
|
||||
for (int i=0; i<COUNT; i++) {
|
||||
Course course = coursesReceived.get(i);
|
||||
Iterable<Mark> marks = course.getMarks();
|
||||
Mark lastMark = null;
|
||||
for (Mark mark : marks) {
|
||||
lastMark = mark;
|
||||
}
|
||||
assertEquals(i, lastMark.getIndex());
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<SailMasterMessage> loadMessages(String raceID) {
|
||||
synchronized (this) {
|
||||
loadMessagesCalled = true;
|
||||
notifyAll();
|
||||
}
|
||||
return messagesToLoad;
|
||||
}
|
||||
}
|
||||
+2
-1
@@ -13,7 +13,8 @@ public interface RaceSpecificMessageLoader {
|
||||
/**
|
||||
* Loads all race-specific messages for the race with ID <code>raceID</code> from the store represented by this
|
||||
* object. The messages can be expected to have their {@link SailMasterMessage#getSequenceNumber() sequence number}
|
||||
* properly set. The result list is expected to be in ascending order regarding this sequence number.
|
||||
* properly set. The result list is expected to be in ascending order regarding this sequence number and always
|
||||
* has to be a valid, non-<code>null</code> and perhaps empty list.
|
||||
*/
|
||||
List<SailMasterMessage> loadMessages(String raceID);
|
||||
}
|
||||
+1
-1
@@ -172,7 +172,7 @@ public class SailMasterConnectorImpl extends SailMasterTransceiverImpl implement
|
||||
raceSpecificMessageBuffers.remove(raceID);
|
||||
}
|
||||
}
|
||||
if (bufferedMessage != null) {
|
||||
if (bufferedMessage != null) { // FIXME compare sequence number to maxSequenceNumber (first create failing test case)
|
||||
notifyListeners(bufferedMessage);
|
||||
}
|
||||
} while (bufferedMessage != null && raceSpecificMessageBuffers.containsKey(raceID));
|
||||
|
||||
Reference in new issue
Block a user