working towards a solution for bug 1910; use asynchronous connectivity for SailMasterConnectorImpl so the timeout can work

This commit is contained in:
Axel Uhl committed 2014-04-29 13:38:42 +02:00
1 parent 4450af2a5f
commit a716539ca4
7 files changed
+129 -107

No files matched your search

@@ -48,7 +48,7 @@ public interface SailMasterConnector {
*/
void addSailMasterListener(SailMasterListener listener) throws UnknownHostException, IOException, InterruptedException;
void removeSailMasterListener(SailMasterListener listener);
void removeSailMasterListener(SailMasterListener listener) throws IOException;
SailMasterMessage receiveMessage(MessageType type) throws InterruptedException;
@@ -90,7 +90,6 @@ public class SailMasterConnectorImpl extends SailMasterTransceiverImpl implement
private final Set<SailMasterListener> listeners;
private final Thread receiverThread;
private boolean stopped;
private boolean connected;
private final String raceId;
private final String raceName;
private final String raceDescription;
@@ -143,45 +142,43 @@ public class SailMasterConnectorImpl extends SailMasterTransceiverImpl implement
this.unprocessedMessagesByType = new HashMap<MessageType, BlockingQueue<SailMasterMessage>>();
receiverThread = new Thread(this, "SwissTiming SailMaster Receiver");
receiverThread.start();
synchronized (this) {
while (!connected) {
wait();
}
}
}
public void run() {
try {
while (!stopped) {
try {
ensureSocketIsOpen();
String receivedMessage = receiveMessage(socket.getInputStream());
if (receivedMessage == null) {
// reached EOF; this means the socket is or can be closed
if (socket != null && !socket.isClosed()) {
socket.close();
}
socket = null;
} else {
SailMasterMessage message = new SailMasterMessageImpl(receivedMessage);
// drop race-specific messages for non-tracked races
if (message.getSequenceNumber() != null) {
maxSequenceNumber = Math.max(maxSequenceNumber, message.getSequenceNumber());
if (maxSequenceNumber <= numberOfStoredMessages) {
notifyListenersStoredDataProgress(raceId, (double) maxSequenceNumber / (double) numberOfStoredMessages);
ensureSocketIsOpen(); // the result of this may be that the connector w\s stopped in between, so another check is required
if (!stopped && socket != null) {
String receivedMessage = receiveMessage(socket.getInputStream());
if (receivedMessage == null) {
// reached EOF; this means the socket is or can be closed
if (socket != null && !socket.isClosed()) {
socket.close();
}
socket = null;
} else {
SailMasterMessage message = new SailMasterMessageImpl(receivedMessage);
// drop race-specific messages for non-tracked races
if (message.getSequenceNumber() != null) {
maxSequenceNumber = Math.max(maxSequenceNumber, message.getSequenceNumber());
if (maxSequenceNumber <= numberOfStoredMessages) {
notifyListenersStoredDataProgress(raceId, (double) maxSequenceNumber
/ (double) numberOfStoredMessages);
}
}
if (message.isResponse()) {
// this is a response for an explicit request
rendevouz(message);
} else if (message.isEvent()) {
// a spontaneous event
logger.fine("notifying message " + message);
notifyListeners(message);
}
if (message.getType() == MessageType._STOPSERVER) {
logger.info("SailMasterConnector received " + MessageType._STOPSERVER.name());
stop();
}
}
if (message.isResponse()) {
// this is a response for an explicit request
rendevouz(message);
} else if (message.isEvent()) {
// a spontaneous event
logger.fine("notifying message " + message);
notifyListeners(message);
}
if (message.getType() == MessageType._STOPSERVER) {
logger.info("SailMasterConnector received " + MessageType._STOPSERVER.name());
stop();
}
}
} catch (SocketException se) {
@@ -205,14 +202,16 @@ public class SailMasterConnectorImpl extends SailMasterTransceiverImpl implement
return result;
}
private synchronized BlockingQueue<SailMasterMessage> getBlockingQueue(MessageType type) {
BlockingQueue<SailMasterMessage> blockingQueue;
blockingQueue = unprocessedMessagesByType.get(type);
if (blockingQueue == null) {
blockingQueue = new LinkedBlockingQueue<SailMasterMessage>();
unprocessedMessagesByType.put(type, blockingQueue);
private BlockingQueue<SailMasterMessage> getBlockingQueue(MessageType type) {
synchronized (unprocessedMessagesByType) {
BlockingQueue<SailMasterMessage> blockingQueue;
blockingQueue = unprocessedMessagesByType.get(type);
if (blockingQueue == null) {
blockingQueue = new LinkedBlockingQueue<SailMasterMessage>();
unprocessedMessagesByType.put(type, blockingQueue);
}
return blockingQueue;
}
return blockingQueue;
}
private void rendevouz(SailMasterMessage message) {
@@ -254,7 +253,7 @@ public class SailMasterConnectorImpl extends SailMasterTransceiverImpl implement
}
private void notifyListenersStoredDataProgress(String raceID, double progress) {
for (SailMasterListener listener : getListeners(raceID)) {
for (SailMasterListener listener : getListeners()) {
try {
listener.storedDataProgress(raceID, progress);
} catch (Exception e) {
@@ -270,7 +269,7 @@ public class SailMasterConnectorImpl extends SailMasterTransceiverImpl implement
int zeroBasedMarkIndex = Integer.valueOf(message.getSections()[2]);
double windDirectionTrueDegrees = Double.valueOf(message.getSections()[3]);
double windSpeedInKnots = Double.valueOf(message.getSections()[4]);
for (SailMasterListener listener : getListeners(message.getRaceID())) {
for (SailMasterListener listener : getListeners()) {
try {
listener.receivedWindData(raceID, zeroBasedMarkIndex, windDirectionTrueDegrees, windSpeedInKnots);
} catch (Exception e) {
@@ -294,7 +293,7 @@ public class SailMasterConnectorImpl extends SailMasterTransceiverImpl implement
parseHHMMSSToMilliseconds(details[2]);
markIndicesRanksAndTimesSinceStartInMilliseconds.add(new Triple<Integer, Integer, Long>(markIndex, rank, timeSinceStartInMilliseconds));
}
for (SailMasterListener listener : getListeners(message.getRaceID())) {
for (SailMasterListener listener : getListeners()) {
try {
listener.receivedTimingData(raceID, boatID, markIndicesRanksAndTimesSinceStartInMilliseconds);
} catch (Exception e) {
@@ -306,7 +305,7 @@ public class SailMasterConnectorImpl extends SailMasterTransceiverImpl implement
private void notifyListenersCAM(SailMasterMessage message) throws ParseException {
List<Triple<Integer, TimePoint, String>> clockAtMarkResults = parseClockAtMarkMessage(message);
for (SailMasterListener listener : getListeners(message.getRaceID())) {
for (SailMasterListener listener : getListeners()) {
try {
listener.receivedClockAtMark(message.getSections()[1], clockAtMarkResults);
} catch (Exception e) {
@@ -318,7 +317,7 @@ public class SailMasterConnectorImpl extends SailMasterTransceiverImpl implement
private void notifyListenersSTL(SailMasterMessage message) {
StartList startListMessage = parseStartListMessage(message);
for (SailMasterListener listener : getListeners(message.getRaceID())) {
for (SailMasterListener listener : getListeners()) {
try {
listener.receivedStartList(message.getSections()[1], startListMessage);
} catch (Exception e) {
@@ -330,7 +329,7 @@ public class SailMasterConnectorImpl extends SailMasterTransceiverImpl implement
private void notifyListenersCCG(SailMasterMessage message) {
Course course = parseCourseConfigurationMessage(message);
for (SailMasterListener listener : getListeners(message.getRaceID())) {
for (SailMasterListener listener : getListeners()) {
try {
listener.receivedCourseConfiguration(message.getSections()[1], course);
} catch (Exception e) {
@@ -342,7 +341,7 @@ public class SailMasterConnectorImpl extends SailMasterTransceiverImpl implement
private void notifyListenersRAC(SailMasterMessage message) {
Iterable<Race> races = parseAvailableRacesMessage(message);
for (SailMasterListener listener : listeners) {
for (SailMasterListener listener : getListeners()) {
try {
listener.receivedAvailableRaces(races);
} catch (Exception e) {
@@ -418,7 +417,7 @@ public class SailMasterConnectorImpl extends SailMasterTransceiverImpl implement
rank, averageSpeedOverGround, velocityMadeGood, distanceToLeader, distanceToNextMark, boatIRM));
}
}
Set<SailMasterListener> allListeners = getListeners(message.getRaceID());
Set<SailMasterListener> allListeners = getListeners();
for (SailMasterListener listener : allListeners) {
try {
listener.receivedRacePositionData(raceID, status, timePoint, startTimeEstimatedStartTime, millisecondsSinceRaceStart,
@@ -430,17 +429,23 @@ public class SailMasterConnectorImpl extends SailMasterTransceiverImpl implement
}
}
private Set<SailMasterListener> getListeners(String raceID) {
return Collections.unmodifiableSet(listeners);
private Set<SailMasterListener> getListeners() {
synchronized (listeners) {
return Collections.unmodifiableSet(listeners);
}
}
@Override
public synchronized void stop() throws IOException {
public void stop() throws IOException {
logger.info("Stopping SailMasterConnector listening on port "+port+" with socket "+socket);
stopped = true;
socket.close();
socket = null;
notifyAll();
if (socket != null) {
socket.close();
socket = null;
}
synchronized (this) {
notifyAll();
}
}
@Override
@@ -450,13 +455,19 @@ public class SailMasterConnectorImpl extends SailMasterTransceiverImpl implement
@Override
public void addSailMasterListener(SailMasterListener listener) throws UnknownHostException, IOException, InterruptedException {
ensureSocketIsOpen();
listeners.add(listener);
synchronized (listeners) {
listeners.add(listener);
}
}
@Override
public void removeSailMasterListener(SailMasterListener listener) {
listeners.remove(listener);
public void removeSailMasterListener(SailMasterListener listener) throws IOException {
synchronized (listeners) {
listeners.remove(listener);
if (listeners.isEmpty()) {
stop();
}
}
}
public SailMasterMessage sendRequestAndGetResponse(MessageType messageType, String... args) throws UnknownHostException, IOException, InterruptedException {
@@ -487,51 +498,52 @@ public class SailMasterConnectorImpl extends SailMasterTransceiverImpl implement
return sailMasterMessage;
}
private synchronized void ensureSocketIsOpen() throws InterruptedException {
while (!stopped && socket == null) {
try {
logger.info("Opening socket to " + host + ":" + port + " and sending " + MessageType.OPN.name() + " message...");
socket = new Socket(host, port);
final OutputStream os = socket.getOutputStream();
final InputStream is = socket.getInputStream();
final SailMasterMessage opnRequest = createSailMasterMessage(MessageType.OPN, raceId);
sendMessage(opnRequest, os);
SailMasterMessage opnResponse = new SailMasterMessageImpl(receiveMessage(is));
if (opnResponse.getType() != MessageType.OPN || !"OK".equals(opnResponse.getSections()[1])) {
logger.info("Recevied non-OK response " + opnResponse + " for our request "
+ opnRequest + ". Closing socket and trying again in 1s...");
closeAndNullSocketAndWaitABit();
} else {
logger.info("Received " + opnResponse + " which seems OK. Continuing with " + MessageType.LSN.name() + " request...");
numberOfStoredMessages = Long.valueOf(opnResponse.getSections()[2]);
List<String> lsnArgs = new ArrayList<>();
lsnArgs.add("ON"); // request live messages always; why not?
if (maxSequenceNumber != -1) {
logger.info("Requesting messages starting from sequence number "+(maxSequenceNumber+1));
// already received a numbered message; ask only for newer messages with greater sequence number
lsnArgs.add(new Long(maxSequenceNumber+1).toString());
}
final SailMasterMessage lsnRequest = createSailMasterMessage(MessageType.LSN, lsnArgs.toArray(new String[0]));
sendMessage(lsnRequest, os);
SailMasterMessageImpl lsnResponse = new SailMasterMessageImpl(receiveMessage(is));
if (lsnResponse.getType() != MessageType.LSN || !"OK".equals(lsnResponse.getSections()[1])) {
logger.info("Received non-OK response " + lsnResponse + " for our request " + lsnRequest
private final Object ensureSocketIsOpenSemaphor = new Object();
private void ensureSocketIsOpen() throws InterruptedException {
synchronized (ensureSocketIsOpenSemaphor) {
while (!stopped && socket == null) {
try {
logger.info("Opening socket to " + host + ":" + port + " and sending " + MessageType.OPN.name()
+ " message...");
socket = new Socket(host, port);
final OutputStream os = socket.getOutputStream();
final InputStream is = socket.getInputStream();
final SailMasterMessage opnRequest = createSailMasterMessage(MessageType.OPN, raceId);
sendMessage(opnRequest, os);
SailMasterMessage opnResponse = new SailMasterMessageImpl(receiveMessage(is));
if (opnResponse.getType() != MessageType.OPN || !"OK".equals(opnResponse.getSections()[1])) {
logger.info("Recevied non-OK response " + opnResponse + " for our request " + opnRequest
+ ". Closing socket and trying again in 1s...");
closeAndNullSocketAndWaitABit();
} else {
logger.info("Received "+lsnResponse+" which seems to be OK. I think we're connected!");
synchronized (this) {
logger.info("...successfully opened socket to " + host + ":" + port);
// TODO now request all messages starting after the last message received so far
connected = true;
notifyAll();
logger.info("Received " + opnResponse + " which seems OK. Continuing with "
+ MessageType.LSN.name() + " request...");
numberOfStoredMessages = Long.valueOf(opnResponse.getSections()[2]);
List<String> lsnArgs = new ArrayList<>();
lsnArgs.add("ON"); // request live messages always; why not?
if (maxSequenceNumber != -1) {
logger.info("Requesting messages starting from sequence number " + (maxSequenceNumber + 1));
// already received a numbered message; ask only for newer messages with greater sequence number
lsnArgs.add(new Long(maxSequenceNumber + 1).toString());
}
final SailMasterMessage lsnRequest = createSailMasterMessage(MessageType.LSN,
lsnArgs.toArray(new String[0]));
sendMessage(lsnRequest, os);
SailMasterMessageImpl lsnResponse = new SailMasterMessageImpl(receiveMessage(is));
if (lsnResponse.getType() != MessageType.LSN || !"OK".equals(lsnResponse.getSections()[1])) {
logger.info("Received non-OK response " + lsnResponse + " for our request " + lsnRequest
+ ". Closing socket and trying again in 1s...");
closeAndNullSocketAndWaitABit();
} else {
logger.info("Received " + lsnResponse
+ " which seems to be OK. I think we're connected to " + host + ":" + port + "!");
}
}
} catch (IOException e) {
logger.log(Level.INFO, "Exception trying to establish SailMaster connection to " + host + ":"
+ port + ". Trying again in 1s.", e);
closeAndNullSocketAndWaitABit();
}
} catch (IOException e) {
logger.log(Level.INFO, "Exception trying to establish SailMaster connection to " + host + ":" + port
+ ". Trying again in 1s.", e);
closeAndNullSocketAndWaitABit();
}
}
}
@@ -185,20 +185,26 @@ public class SwissTimingRaceTrackerImpl extends AbstractRaceTrackerImpl implemen
public Set<RaceDefinition> getRaces(long timeoutInMilliseconds) {
long start = System.currentTimeMillis();
synchronized (this) {
RaceDefinition result = race;
RaceDefinition preResult = race;
boolean interrupted = false;
while ((System.currentTimeMillis()-start < timeoutInMilliseconds) && !interrupted && result == null) {
while ((System.currentTimeMillis()-start < timeoutInMilliseconds) && !interrupted && preResult == null) {
try {
long timeToWait = timeoutInMilliseconds - (System.currentTimeMillis() - start);
if (timeToWait > 0) {
this.wait(timeToWait);
}
result = race;
preResult = race;
} catch (InterruptedException e) {
interrupted = true;
}
}
return result == null ? null : Collections.singleton(result);
final Set<RaceDefinition> result;
if (preResult == null) {
result = Collections.emptySet();
} else {
result = Collections.singleton(preResult);
}
return result;
}
}
@@ -241,10 +241,10 @@ public interface DomainFactory {
* <code>timeoutInMilliseconds</code> milliseconds have passed and the race definition is found not to have shown up
* until then, <code>null</code> is returned. The unblocking may be deferred even beyond
* <code>timeoutInMilliseconds</code> in case no modifications happen on the set of races cached by this factory.
* @param raceId TODO
*
* @param timeoutInMilliseconds
* passing -1 means an infinite timeout; 0 means return immediately with <code>null</code> as result if no
* race definition is found for <code>race</code>.
* passing -1 means an infinite timeout; 0 means return immediately with <code>null</code> as result if
* no race definition is found for <code>race</code>.
*/
RaceDefinition getAndWaitForRaceDefinition(UUID raceId, long timeoutInMilliseconds);
@@ -52,7 +52,10 @@ public class RaceHandleImpl implements RacesHandle {
public Set<RaceDefinition> getRaces(long timeoutInMilliseconds) {
Set<RaceDefinition> result = new HashSet<RaceDefinition>();
for (Race race : tractracEvent.getRaceList()) {
result.add(domainFactory.getAndWaitForRaceDefinition(race.getId(), timeoutInMilliseconds));
final RaceDefinition raceDefinition = domainFactory.getAndWaitForRaceDefinition(race.getId(), timeoutInMilliseconds);
if (raceDefinition != null) { // may have time-outed
result.add(raceDefinition);
}
}
return result;
}
@@ -27,7 +27,7 @@ public interface RacesHandle {
* Fetch the race definition. If the race definition represented by this handle hasn't been created yet, the call
* blocks until such a definition is provided by another call, usually by the {@link RaceCourseReceiver}. If
* <code>timeoutInMilliseconds</code> milliseconds have passed and the race definition is found not to have
* shown up until then, <code>null</code> is returned. The unblocking may be deferred even beyond
* shown up until then, a valid but empty set is returned. The unblocking may be deferred even beyond
* <code>timeoutInMilliseconds</code> in case no modifications happen on the {@link DomainFactory}'s
* set of races during that time.
*/
@@ -1280,6 +1280,7 @@ public class RacingEventServiceImpl implements RacingEventServiceWithTestSupport
ScheduledFuture<?> task = getScheduler().schedule(new Runnable() {
@Override
public void run() {
if (tracker.getRaces() == null || tracker.getRaces().isEmpty()) {
try {
Regatta regatta = tracker.getRegatta();