From fd9422dc582817b7e1ba2a9e0dd3b5fa3ebd52d9 Mon Sep 17 00:00:00 2001 From: Axel Uhl Date: Tue, 14 Feb 2012 17:29:28 +0100 Subject: [PATCH] attempt to factor out the HTTP-based send/receive streaming aspect; but the request input stream closes fast... --- .../sailing/domain/common}/Base64Utils.java | 2 +- .../ui/usermanagement/UserManagementPage.java | 2 +- .../ExpeditionThroughHttpPostServletTest.java | 40 ++--- .../server/ExpeditionHttpReceiver.java | 142 ++++++++++++++++++ .../server/impl/AbstractHttpPostServlet.java | 126 ++++++++++++++++ .../ExpeditionThroughHttpPostServlet.java | 98 ++++-------- 6 files changed, 310 insertions(+), 100 deletions(-) rename java/{com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/shared => com.sap.sailing.domain.common/src/com/sap/sailing/domain/common}/Base64Utils.java (96%) create mode 100755 java/com.sap.sailing.server/src/com/sap/sailing/server/ExpeditionHttpReceiver.java create mode 100755 java/com.sap.sailing.server/src/com/sap/sailing/server/impl/AbstractHttpPostServlet.java diff --git a/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/shared/Base64Utils.java b/java/com.sap.sailing.domain.common/src/com/sap/sailing/domain/common/Base64Utils.java similarity index 96% rename from java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/shared/Base64Utils.java rename to java/com.sap.sailing.domain.common/src/com/sap/sailing/domain/common/Base64Utils.java index ef15c865581..118c9016173 100755 --- a/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/shared/Base64Utils.java +++ b/java/com.sap.sailing.domain.common/src/com/sap/sailing/domain/common/Base64Utils.java @@ -1,4 +1,4 @@ -package com.sap.sailing.gwt.ui.shared; +package com.sap.sailing.domain.common; /* * Copyright 2009 Google Inc. diff --git a/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/usermanagement/UserManagementPage.java b/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/usermanagement/UserManagementPage.java index aac38e9ddcd..459d261fb02 100755 --- a/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/usermanagement/UserManagementPage.java +++ b/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/usermanagement/UserManagementPage.java @@ -11,8 +11,8 @@ import com.google.gwt.user.client.ui.PasswordTextBox; import com.google.gwt.user.client.ui.RootPanel; import com.google.gwt.user.client.ui.TextBox; import com.google.gwt.user.client.ui.VerticalPanel; +import com.sap.sailing.domain.common.Base64Utils; import com.sap.sailing.gwt.ui.client.AbstractEntryPoint; -import com.sap.sailing.gwt.ui.shared.Base64Utils; public class UserManagementPage extends AbstractEntryPoint { @Override diff --git a/java/com.sap.sailing.server.test/src/com/sap/sailing/server/test/ExpeditionThroughHttpPostServletTest.java b/java/com.sap.sailing.server.test/src/com/sap/sailing/server/test/ExpeditionThroughHttpPostServletTest.java index 1e5eac35718..3a687a2bf4c 100755 --- a/java/com.sap.sailing.server.test/src/com/sap/sailing/server/test/ExpeditionThroughHttpPostServletTest.java +++ b/java/com.sap.sailing.server.test/src/com/sap/sailing/server/test/ExpeditionThroughHttpPostServletTest.java @@ -1,45 +1,25 @@ package com.sap.sailing.server.test; -import java.io.BufferedReader; import java.io.IOException; -import java.io.InputStreamReader; -import java.io.PrintWriter; -import java.net.HttpURLConnection; import java.net.URL; -import org.junit.Ignore; import org.junit.Test; +import com.sap.sailing.server.ExpeditionHttpReceiver; + public class ExpeditionThroughHttpPostServletTest { - @Ignore // To run, the OSGi-based Jetty needs to run and listen on port 8888" +// @Ignore // To run, the OSGi-based Jetty needs to run and listen on port 8888" @Test public void testConnectDisconnect() throws IOException, InterruptedException { int jettyPort = 8888; URL url = new URL("http://localhost:"+jettyPort+"/sailingserver/expedition"); - final HttpURLConnection connection = (HttpURLConnection) url.openConnection(); - connection.setChunkedStreamingMode(/* chunklen */ 8192); - connection.setDoOutput(true); - connection.setRequestMethod("POST"); - PrintWriter requestWriter = new PrintWriter(connection.getOutputStream()); - requestWriter.println(""); - requestWriter.flush(); - Thread reader = new Thread(new Runnable() { - public void run() { - try { - BufferedReader br = new BufferedReader(new InputStreamReader(connection.getInputStream())); - String line = br.readLine(); - while (line != null) { - System.out.println(line); - line = br.readLine(); - } - br.close(); - } catch (IOException e) { - e.printStackTrace(); - } + ExpeditionHttpReceiver httpReceiver = new ExpeditionHttpReceiver(url, new ExpeditionHttpReceiver.Receiver() { + @Override + public boolean received(byte[] bytes) { + System.out.println(new String(bytes)); + return false; // don't stop } - }, "ExpeditionThroughHttpPostServletTest reader"); - reader.start(); - reader.join(); - requestWriter.close(); + }); + httpReceiver.connect(); } } diff --git a/java/com.sap.sailing.server/src/com/sap/sailing/server/ExpeditionHttpReceiver.java b/java/com.sap.sailing.server/src/com/sap/sailing/server/ExpeditionHttpReceiver.java new file mode 100755 index 00000000000..7ab8c1f8a35 --- /dev/null +++ b/java/com.sap.sailing.server/src/com/sap/sailing/server/ExpeditionHttpReceiver.java @@ -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 true if the receiver shall terminate and close the connection, false + * 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 "<ping>" 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); + } + } +} diff --git a/java/com.sap.sailing.server/src/com/sap/sailing/server/impl/AbstractHttpPostServlet.java b/java/com.sap.sailing.server/src/com/sap/sailing/server/impl/AbstractHttpPostServlet.java new file mode 100755 index 00000000000..50d31588b9a --- /dev/null +++ b/java/com.sap.sailing.server/src/com/sap/sailing/server/impl/AbstractHttpPostServlet.java @@ -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 + * "" terminated with a newline character through the request stream, and this servlet will respond with the + * message "" 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 = ""; + public static final String 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 writer. To stop the forwarding process, + * call the runToStop 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 "\n" message comes along, a "\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(); + } + } + } + +} diff --git a/java/com.sap.sailing.server/src/com/sap/sailing/server/impl/ExpeditionThroughHttpPostServlet.java b/java/com.sap.sailing.server/src/com/sap/sailing/server/impl/ExpeditionThroughHttpPostServlet.java index 54eb5859c18..ba354177cd7 100755 --- a/java/com.sap.sailing.server/src/com/sap/sailing/server/impl/ExpeditionThroughHttpPostServlet.java +++ b/java/com.sap.sailing.server/src/com/sap/sailing/server/impl/ExpeditionThroughHttpPostServlet.java @@ -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 "" terminated with a newline character through the request stream, + * and this servlet will respond with the string "" 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 writer. To stop the forwarding process, + * call the runToStop 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 "\n" message comes along, a "\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("")) { - synchronized (responseWriter) { - responseWriter.println(""); - responseWriter.flush(); - } - } - line = br.readLine(); - } - } catch (IOException e) { - e.printStackTrace(); - } - } } }