use a send-only heartbeat for Expedition through HTTP connection

This commit is contained in:
Axel Uhl committed 2012-02-14 18:40:53 +01:00
1 parent fd9422dc58
commit 2c7656eeda
5 files changed
+135 -107

No files matched your search

@@ -21,5 +21,7 @@ public class ExpeditionThroughHttpPostServletTest {
}
});
httpReceiver.connect();
Thread.sleep(100000);
httpReceiver.stop();
}
}
@@ -3,10 +3,12 @@ package com.sap.sailing.server;
import java.io.IOException;
import java.io.InputStream;
import java.io.InputStreamReader;
import java.io.PrintWriter;
import java.io.Reader;
import java.net.HttpURLConnection;
import java.net.URL;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.logging.Logger;
import com.sap.sailing.domain.common.Base64Utils;
@@ -20,7 +22,7 @@ import com.sap.sailing.server.impl.AbstractHttpPostServlet;
* @author Axel Uhl (D043530)
*
*/
public class ExpeditionHttpReceiver implements Runnable {
public class ExpeditionHttpReceiver {
public interface Receiver {
/**
* Called when the receiver has received some bytes from the remote end.
@@ -34,11 +36,16 @@ public class ExpeditionHttpReceiver implements Runnable {
}
private static final int BUF_SIZE = 1<<16;
private static final long HEARTBEAT_TIME_IN_MILLISECONDS = 1000;
private static final Logger logger = Logger.getLogger(ExpeditionHttpReceiver.class.getName());
private final Receiver receiver;
private final URL url;
private PrintWriter requestWriter;
private long timestampOfLastHeartbeatReceived;
/**
* The input stream can be used to unblock a read in order to terminate the receiver.
*/
private InputStream inputStream;
private boolean stop = false;
public ExpeditionHttpReceiver(URL url, Receiver receiver) {
@@ -47,44 +54,48 @@ public class ExpeditionHttpReceiver implements Runnable {
}
public void connect() throws IOException, InterruptedException {
final HttpURLConnection connection = (HttpURLConnection) url.openConnection();
connection.setChunkedStreamingMode(/* chunklen */ BUF_SIZE);
connection.setDoOutput(true);
connection.setRequestMethod("POST");
connection.connect();
requestWriter = new PrintWriter(connection.getOutputStream());
requestWriter.println(AbstractHttpPostServlet.PING); // ensure the stream writes through to the other end
requestWriter.flush();
final InputStream inputStream = connection.getInputStream();
new Thread(this, getClass().getName()+" Heartbeat").start();
establishConnection(); // performs the actual HTTP connect request, connecting to the servlet
Thread reader = new Thread(new Runnable() {
public void run() {
StringBuilder bos = new StringBuilder();
try {
Reader reader = new InputStreamReader(inputStream);
char[] buf = new char[BUF_SIZE];
int read = reader.read(buf);
while (!stop && read != -1) {
for (int i = 0; i < read; i++) {
if (buf[i] == 0) {
// message terminator; one message received
stop = receivedMessage(bos.toString());
bos.delete(0, bos.length());
} else {
bos.append(buf[i]);
while (!stop) {
try {
Reader reader = new InputStreamReader(inputStream);
char[] buf = new char[BUF_SIZE];
int read = reader.read(buf);
while (!stop && read != -1) {
for (int i = 0; i < read; i++) {
if (buf[i] == 0) {
// message terminator; one message received
stop = receivedMessage(bos.toString());
bos.delete(0, bos.length());
} else {
bos.append(buf[i]);
}
}
if (!stop) {
read = reader.read(buf);
}
}
if (!stop) {
read = reader.read(buf);
logger.info("Reached EOF");
reader.close();
} catch (IOException e) {
logger.throwing(ExpeditionHttpReceiver.class.getName(), "connect", e);
}
if (!stop) {
logger.info("Reconnecting because not stopped");
try {
establishConnection();
} catch (IOException e) {
logger.info("Can't re-connect. Giving up.");
logger.throwing(ExpeditionHttpReceiver.class.getName(), "connect", e);
stop = true;
}
}
reader.close();
} catch (IOException e) {
logger.throwing(ExpeditionHttpReceiver.class.getName(), "connect", e);
}
}
protected boolean receivedMessage(String bos) {
private boolean receivedMessage(String bos) {
boolean stopReceiving = false;
if (bos.equals(AbstractHttpPostServlet.PONG)) {
receivedHeartbeatResponse();
@@ -96,47 +107,55 @@ public class ExpeditionHttpReceiver implements Runnable {
}, getClass().getName()+" reader");
reader.start();
reader.join();
stop = true; // make sure the heartbeat thread stops
if (requestWriter != null) {
requestWriter.close();
}
scheduleTimeoutHandler();
}
private void scheduleTimeoutHandler() {
final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(1);
Runnable timeoutChecker = new Runnable() {
@Override
public void run() {
if (System.currentTimeMillis() - timestampOfLastHeartbeatReceived > 5*AbstractHttpPostServlet.HEARTBEAT_TIME_IN_MILLISECONDS) {
// TIMEOUT; abort
logger.info("Timeout. Didn't receive a heartbeat through my HTTP connection for "+
(5*AbstractHttpPostServlet.HEARTBEAT_TIME_IN_MILLISECONDS)+"ms");
try {
stop();
} catch (IOException e) {
logger.throwing(ExpeditionHttpReceiver.class.getName(), "scheduleTimeoutHandler", e);
}
} else {
scheduler.schedule(this, AbstractHttpPostServlet.HEARTBEAT_TIME_IN_MILLISECONDS, TimeUnit.MILLISECONDS);
}
}
};
scheduler.schedule(timeoutChecker, AbstractHttpPostServlet.HEARTBEAT_TIME_IN_MILLISECONDS, TimeUnit.MILLISECONDS);
}
/**
* connects to the remote servlet using {@link #url} and binds the response input stream to {@link #inputStream}
*/
private void establishConnection() throws IOException {
HttpURLConnection connection = (HttpURLConnection) url.openConnection();
inputStream = connection.getInputStream();
}
private void receivedHeartbeatResponse() {
// TODO do we want to keep track of received heartbeat responses?
logger.finest("received expedition HTTP heartbeat");
timestampOfLastHeartbeatReceived = System.currentTimeMillis();
}
/**
* Stops the receiver by closing the writer on the request stream. This will let the server read EOF on the
* request stream which causes the server to also terminate the sending of data, closing its response stream.
* This in turn will let the reader return with an EOF.
* @throws IOException
*/
public void stop() {
public void stop() throws IOException {
logger.info("Stopping expedition HTTP receiver");
stop = true;
if (requestWriter != null) {
requestWriter.close();
}
}
/**
* Implements the heartbeat by sending a "&lt;ping&gt;" message every {@link #HEARTBEAT_TIME_IN_MILLISECONDS} milliseconds.
* If that fails with an exception, the {@link #stop} method is called.
*/
@Override
public void run() {
try {
while (!stop && requestWriter != null) {
Thread.sleep(HEARTBEAT_TIME_IN_MILLISECONDS);
if (requestWriter != null) {
synchronized (requestWriter) {
requestWriter.println(AbstractHttpPostServlet.PING);
requestWriter.flush();
}
}
}
} catch (InterruptedException e) {
logger.throwing(ExpeditionHttpReceiver.class.getName(), "run", e);
if (inputStream != null) {
inputStream.close();
}
}
}
@@ -1,16 +1,12 @@
package com.sap.sailing.server.impl;
import java.io.BufferedReader;
import java.io.IOException;
import java.io.InputStream;
import java.io.InputStreamReader;
import java.io.PrintWriter;
import java.io.Writer;
import java.net.SocketException;
import java.util.logging.Logger;
import javax.servlet.ServletException;
import javax.servlet.ServletInputStream;
import javax.servlet.http.HttpServletRequest;
import javax.servlet.http.HttpServletResponse;
@@ -29,36 +25,42 @@ import com.sap.sailing.server.Servlet;
*
*/
public abstract class AbstractHttpPostServlet extends Servlet {
public static final String PING = "<ping>";
/**
* Every so many milliseconds the servlet will send a {@link #PONG} message. This shows clients that the connection
* is still alive. If after several such intervals the client hasn't received a {@link #PONG} message it seems reasonable
* to assume the connection has died.
*/
public static final long HEARTBEAT_TIME_IN_MILLISECONDS = 5000;
public static final String PONG = "<pong>";
private static final Logger logger = Logger.getLogger(AbstractHttpPostServlet.class.getName());
private static final long timeoutInMilliseconds = 60000;
private long timeInMillisOfLastExpeditionMessageReceived;
private boolean stop;
private static final long serialVersionUID = 6034769972654796465L;
@Override
protected void doPost(HttpServletRequest req, HttpServletResponse resp) throws ServletException, IOException {
protected void doGet(HttpServletRequest req, HttpServletResponse resp) throws ServletException, IOException {
final PrintWriter writer = resp.getWriter();
final ServletInputStream requestInputStream = req.getInputStream();
Thread heartbeatHandler = new Thread(new HeartbeatHandler(requestInputStream, writer), getClass().getName()
HeartbeatHandler heartbeat = new HeartbeatHandler(writer);
Thread heartbeatHandler = new Thread(heartbeat, getClass().getName()
+ " HeartbeatHandler " + Thread.currentThread().getId());
startSendingResponse(writer, new Runnable() {
public void run() {
try {
requestInputStream.close();
} catch (IOException e) {
logger.throwing(ExpeditionThroughHttpPostServlet.class.getName(), "run", e);
}
stop = true;
}
});
heartbeatHandler.start();
try {
Thread.sleep(timeoutInMilliseconds);
while (System.currentTimeMillis()-timeInMillisOfLastExpeditionMessageReceived < timeoutInMilliseconds) {
while (!stop && System.currentTimeMillis()-timeInMillisOfLastExpeditionMessageReceived < timeoutInMilliseconds) {
Thread.sleep(timeInMillisOfLastExpeditionMessageReceived + timeoutInMilliseconds - System.currentTimeMillis());
}
if (stop) {
logger.info(getClass().getName()+" was explicitly stopped, e.g., because client closed connection");
}
logger.info("Terminating "+getClass().getName()+" doPost after not receiving anything for "+timeoutInMilliseconds+"ms");
heartbeat.stop();
heartbeatHandler.join();
} catch (InterruptedException e) {
e.printStackTrace();
@@ -72,12 +74,21 @@ public abstract class AbstractHttpPostServlet extends Servlet {
* @param bytes
* @throws IOException
*/
protected void send(Writer writer, byte[] message) throws IOException {
protected void send(Writer writer, byte[] message) {
String bytes = Base64Utils.toBase64(message);
sendString(writer, bytes);
}
protected void sendString(Writer writer, String s) {
timeInMillisOfLastExpeditionMessageReceived = System.currentTimeMillis();
synchronized (writer) {
writer.write(Base64Utils.toBase64(message));
writer.write(0); // terminate message with 0 character
writer.flush();
try {
synchronized (writer) {
writer.write(s);
writer.write(0); // terminate message with 0 character
writer.flush();
}
} catch (IOException e) {
stop = true; // probably the client closed the connection
}
}
@@ -89,36 +100,36 @@ public abstract class AbstractHttpPostServlet extends Servlet {
abstract protected void startSendingResponse(final Writer writer, final Runnable runToStop) throws SocketException;
/**
* Reads from the request input stream. If a "<ping>\n" message comes along, a "<pong>\n" message will be sent back.
* This will allow a client to regularly check live-ness of the connection and reconnect if needed.
* Every {@link #HEARTBEAT_TIME_IN_MILLISECONDS} milliseconds, a "<pong>\n" message will be sent to {@link #responseWriter}.
* This will allow a client to regularly check live-ness of the connection and reconnect if needed. If {@link #stop} is
* called on this object, the heartbeat sending is stopped and the {@link #run} method returns.
*
* @author Axel Uhl (D043530)
*/
private class HeartbeatHandler implements Runnable {
private final InputStream requestInputStream;
private boolean stop = false;
private final PrintWriter responseWriter;
public HeartbeatHandler(InputStream requestInputStream, PrintWriter responseWriter) {
public HeartbeatHandler(PrintWriter responseWriter) {
super();
this.requestInputStream = requestInputStream;
this.responseWriter = responseWriter;
}
public void stop() {
stop = true;
}
@Override
public void run() {
BufferedReader br = new BufferedReader(new InputStreamReader(requestInputStream));
try {
String line = br.readLine();
while (line != null) {
if (line.equals(PING)) {
synchronized (responseWriter) {
send(responseWriter, PONG.getBytes());
}
}
line = br.readLine();
while (!stop) {
sendString(responseWriter, PONG);
Thread.sleep(HEARTBEAT_TIME_IN_MILLISECONDS);
}
} catch (IOException e) {
e.printStackTrace();
} catch (Exception e) {
logger.info("Terminating heartbeat on "+AbstractHttpPostServlet.this.getClass().getName()+
" because of exception "+e);
logger.throwing(getClass().getName(), "run", e);
}
}
}
@@ -1,9 +1,7 @@
package com.sap.sailing.server.impl;
import java.io.IOException;
import java.io.Writer;
import java.net.SocketException;
import java.util.logging.Logger;
import com.sap.sailing.expeditionconnector.ExpeditionListener;
import com.sap.sailing.expeditionconnector.ExpeditionMessage;
@@ -23,7 +21,6 @@ import com.sap.sailing.server.RacingEventService;
*/
public class ExpeditionThroughHttpPostServlet extends AbstractHttpPostServlet {
private static final long serialVersionUID = 4409173886816756920L;
private static final Logger logger = Logger.getLogger(ExpeditionThroughHttpPostServlet.class.getName());
/**
* Used to start sending the response. This may well happen in a separate thread spawned by this method or, e.g., by
@@ -36,12 +33,7 @@ public class ExpeditionThroughHttpPostServlet extends AbstractHttpPostServlet {
@Override
public void received(ExpeditionMessage message) {
synchronized (writer) {
try {
send(writer, message.getOriginalMessage().getBytes());
} catch (IOException e) {
logger.throwing(ExpeditionThroughHttpPostServlet.class.getName(), "received", e);
runToStop.run();
}
send(writer, message.getOriginalMessage().getBytes());
}
}
}, /* validMessagesOnly */ false);
@@ -434,7 +434,7 @@ public class RacingEventServiceImpl implements RacingEventService, EventFetcher,
* when the tracker is stopped or has successfully received the race
*/
private ScheduledFuture<?> scheduleAbortTrackerAfterInitialTimeout(final RaceTracker tracker, final long timeoutInMilliseconds) {
ScheduledFuture<?> task = scheduler.schedule(new Runnable() {
ScheduledFuture<?> task = getScheduler().schedule(new Runnable() {
@Override public void run() {
if (tracker.getRaces() == null || tracker.getRaces().isEmpty()) {
try {
@@ -747,4 +747,8 @@ public class RacingEventServiceImpl implements RacingEventService, EventFetcher,
receiver.removeListener(listener);
}
private ScheduledExecutorService getScheduler() {
return scheduler;
}
}