attempt to factor out the HTTP-based send/receive streaming aspect; but the request input stream closes fast...

This commit is contained in:
Axel Uhl committed 2012-02-14 17:29:28 +01:00
1 parent 099418ca21
commit fd9422dc58
6 files changed
+310 -100

No files matched your search

@@ -0,0 +1,142 @@
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.logging.Logger;
import com.sap.sailing.domain.common.Base64Utils;
import com.sap.sailing.server.impl.AbstractHttpPostServlet;
/**
* Receives data from a remote servlet, trying to keep the connection open until {@link #stop stopped}. After constructing
* an instance, clients need to call {@link #connect()} to actually start the process of receiving data through the HTTP
* connection.
*
* @author Axel Uhl (D043530)
*
*/
public class ExpeditionHttpReceiver implements Runnable {
public interface Receiver {
/**
* Called when the receiver has received some bytes from the remote end.
* @param bytes
* 0..count-1 hold the bytes received
*
* @return <code>true</code> if the receiver shall terminate and close the connection, <code>false</code>
* otherwise
*/
boolean received(byte[] bytes);
}
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 boolean stop = false;
public ExpeditionHttpReceiver(URL url, Receiver receiver) {
this.receiver = receiver;
this.url = url;
}
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();
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]);
}
}
if (!stop) {
read = reader.read(buf);
}
}
reader.close();
} catch (IOException e) {
logger.throwing(ExpeditionHttpReceiver.class.getName(), "connect", e);
}
}
protected boolean receivedMessage(String bos) {
boolean stopReceiving = false;
if (bos.equals(AbstractHttpPostServlet.PONG)) {
receivedHeartbeatResponse();
} else {
stopReceiving = receiver.received(Base64Utils.fromBase64(bos));
}
return stopReceiving;
}
}, getClass().getName()+" reader");
reader.start();
reader.join();
stop = true; // make sure the heartbeat thread stops
if (requestWriter != null) {
requestWriter.close();
}
}
private void receivedHeartbeatResponse() {
// TODO do we want to keep track of received heartbeat responses?
}
/**
* 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.
*/
public void stop() {
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);
}
}
}
@@ -0,0 +1,126 @@
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;
import com.sap.sailing.domain.common.Base64Utils;
import com.sap.sailing.server.Servlet;
/**
* Subclasses can be used to retrieve messages through an HTTP connection. The connection remains open until the client
* closes the request stream. Of course, network errors can occur, so the connection may be closed for other reasons,
* too. The client may not necessarily become aware of the connection breaking. Therefore, clients can send the string
* "<ping>" terminated with a newline character through the request stream, and this servlet will respond with the
* message "<pong>" on the output stream. The ping/pong messages interleave the message stream but never cut a single
* message into several parts.
*
* @author Axel Uhl (D043530)
*
*/
public abstract class AbstractHttpPostServlet extends Servlet {
public static final String PING = "<ping>";
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 static final long serialVersionUID = 6034769972654796465L;
@Override
protected void doPost(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 " + Thread.currentThread().getId());
startSendingResponse(writer, new Runnable() {
public void run() {
try {
requestInputStream.close();
} catch (IOException e) {
logger.throwing(ExpeditionThroughHttpPostServlet.class.getName(), "run", e);
}
}
});
heartbeatHandler.start();
try {
Thread.sleep(timeoutInMilliseconds);
while (System.currentTimeMillis()-timeInMillisOfLastExpeditionMessageReceived < timeoutInMilliseconds) {
Thread.sleep(timeInMillisOfLastExpeditionMessageReceived + timeoutInMilliseconds - System.currentTimeMillis());
}
logger.info("Terminating "+getClass().getName()+" doPost after not receiving anything for "+timeoutInMilliseconds+"ms");
heartbeatHandler.join();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
/**
* Sends a single message across to the receiving client. To do so, the message is Base64 encoded before being sent
* as a string. The message is terminated with a 0 byte.
* @param writer
* @param bytes
* @throws IOException
*/
protected void send(Writer writer, byte[] message) throws IOException {
timeInMillisOfLastExpeditionMessageReceived = System.currentTimeMillis();
synchronized (writer) {
writer.write(Base64Utils.toBase64(message));
writer.write(0); // terminate message with 0 character
writer.flush();
}
}
/**
* Used to start sending the response. This may well happen in a separate thread spawned by this method or, e.g., by
* registering for receiving data and sending it to the <code>writer</code>. To stop the forwarding process,
* call the <code>runToStop</code> object's {@link Runnable#run()} method.
*/
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.
*
* @author Axel Uhl (D043530)
*/
private class HeartbeatHandler implements Runnable {
private final InputStream requestInputStream;
private final PrintWriter responseWriter;
public HeartbeatHandler(InputStream requestInputStream, PrintWriter responseWriter) {
super();
this.requestInputStream = requestInputStream;
this.responseWriter = responseWriter;
}
@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();
}
} catch (IOException e) {
e.printStackTrace();
}
}
}
}
@@ -1,87 +1,49 @@
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 javax.servlet.ServletException;
import javax.servlet.http.HttpServletRequest;
import javax.servlet.http.HttpServletResponse;
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;
import com.sap.sailing.server.RacingEventService;
import com.sap.sailing.server.Servlet;
public class ExpeditionThroughHttpPostServlet extends Servlet {
private static final long timeoutInMilliseconds = 60000;
private long timeInMillisOfLastExpeditionMessageReceived;
/**
* Clients can use this servlet to retrieve Expedition messages through an HTTP connection. The connection remains
* open until the client closes the request stream. Of course, network errors can occur, so the connection may
* be closed for other reasons, too. The client may not necessarily become aware of the connection breaking.
* Therefore, clients can send the string "<ping>" terminated with a newline character through the request stream,
* and this servlet will respond with the string "<pong>" terminated by a newline on the output stream. The
* ping/pong messages interleave a stream of Expedition messages but never cut a single expedition message into
* several pieces.
*
* @author Axel Uhl (D043530)
*
*/
public class ExpeditionThroughHttpPostServlet extends AbstractHttpPostServlet {
private static final long serialVersionUID = 4409173886816756920L;
private static final Logger logger = Logger.getLogger(ExpeditionThroughHttpPostServlet.class.getName());
private static final long serialVersionUID = 6034769972654796465L;
@Override
protected void doPost(HttpServletRequest req, HttpServletResponse resp) throws ServletException, IOException {
/**
* Used to start sending the response. This may well happen in a separate thread spawned by this method or, e.g., by
* registering for receiving data and sending it to the <code>writer</code>. To stop the forwarding process,
* call the <code>runToStop</code> object's {@link Runnable#run()} method.
*/
protected void startSendingResponse(final Writer writer, final Runnable runToStop) throws SocketException {
RacingEventService service = getService();
final PrintWriter writer = resp.getWriter();
Thread pingPongHandler = new Thread(new PingPongHandler(req.getInputStream(), writer), getClass().getName()
+ " PingPongHandler " + Thread.currentThread().getId());
service.addExpeditionListener(new ExpeditionListener() {
@Override
public void received(ExpeditionMessage message) {
timeInMillisOfLastExpeditionMessageReceived = System.currentTimeMillis();
synchronized (writer) {
writer.println(message.getOriginalMessage());
writer.flush();
try {
send(writer, message.getOriginalMessage().getBytes());
} catch (IOException e) {
logger.throwing(ExpeditionThroughHttpPostServlet.class.getName(), "received", e);
runToStop.run();
}
}
}
}, /* validMessagesOnly */ false);
pingPongHandler.start();
try {
Thread.sleep(timeoutInMilliseconds);
while (System.currentTimeMillis()-timeInMillisOfLastExpeditionMessageReceived < timeoutInMilliseconds) {
Thread.sleep(timeInMillisOfLastExpeditionMessageReceived + timeoutInMilliseconds - System.currentTimeMillis());
}
pingPongHandler.join();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
/**
* 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.
*
* @author Axel Uhl (D043530)
*/
private class PingPongHandler implements Runnable {
private final InputStream requestInputStream;
private final PrintWriter responseWriter;
public PingPongHandler(InputStream requestInputStream, PrintWriter responseWriter) {
super();
this.requestInputStream = requestInputStream;
this.responseWriter = responseWriter;
}
@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) {
responseWriter.println("<pong>");
responseWriter.flush();
}
}
line = br.readLine();
}
} catch (IOException e) {
e.printStackTrace();
}
}
}
}