Revision: 6880 Author: [email protected] Date: Thu Nov 12 13:26:49 2009 Log: Adds a start() method to MessageTransport. The message transport threads are no longer started in the constructor; they are only started once the start() method is called.
Review by: mmendez http://code.google.com/p/google-web-toolkit/source/detail?r=6880 Modified: /trunk/dev/core/src/com/google/gwt/dev/shell/remoteui/MessageTransport.java /trunk/dev/core/src/com/google/gwt/dev/shell/remoteui/RemoteUI.java /trunk/dev/core/test/com/google/gwt/dev/shell/remoteui/MessageTransportTest.java ======================================= --- /trunk/dev/core/src/com/google/gwt/dev/shell/remoteui/MessageTransport.java Wed Nov 11 18:23:59 2009 +++ /trunk/dev/core/src/com/google/gwt/dev/shell/remoteui/MessageTransport.java Thu Nov 12 13:26:49 2009 @@ -33,6 +33,7 @@ import java.util.concurrent.Executors; import java.util.concurrent.Future; import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.locks.Condition; import java.util.concurrent.locks.Lock; @@ -228,18 +229,18 @@ private static final int DEFAULT_SERVICE_THREADS = 2; - private final Thread messageProcessingThread; + private final AtomicBoolean isStarted = new AtomicBoolean(false); private final AtomicInteger nextMessageId = new AtomicInteger(); private final RequestProcessor requestProcessor; private final LinkedBlockingQueue<PendingSend> sendQueue = new LinkedBlockingQueue<PendingSend>(); - private final Thread sendThread; private final ExecutorService serverRequestExecutor; private final PendingRequestMap pendingRequestMap = new PendingRequestMap(); private final TerminationCallback terminationCallback; + private final InputStream inputStream; + private final OutputStream outputStream; /** * Create a new instance using the given streams and request processor. - * Closing either stream will cause the termination of the transport. * * @param inputStream an input stream for reading messages * @param outputStream an output stream for writing messages @@ -253,10 +254,53 @@ TerminationCallback terminationCallback) { this.requestProcessor = requestProcessor; this.terminationCallback = terminationCallback; + this.inputStream = inputStream; + this.outputStream = outputStream; serverRequestExecutor = Executors.newFixedThreadPool(DEFAULT_SERVICE_THREADS); + } + + /** + * Asynchronously executes the request on a remote server. + * + * @param requestMessage The request to execute + * + * @return a {...@link Future} that can be used to access the server's response + */ + public Future<Response> executeRequestAsync(final Request requestMessage) { + Future<Response> responseFuture = serverRequestExecutor.submit(new Callable<Response>() { + public Response call() throws Exception { + Message.Builder messageBuilder = Message.newBuilder(); + int messageId = nextMessageId.getAndIncrement(); + messageBuilder.setMessageId(messageId); + messageBuilder.setMessageType(Message.MessageType.REQUEST); + messageBuilder.setRequest(requestMessage); + + Message message = messageBuilder.build(); + PendingRequest pendingRequest = new PendingRequest(message); + sendQueue.put(pendingRequest); + + return pendingRequest.waitForResponse(); + } + }); + + return responseFuture; + } + + /** + * Starts up the message transport. The message transport creates its own + * threads, so it is not necessary to invoke this method from a separate + * thread. + * + * Closing either stream will cause the termination of the transport. + */ + public void start() { + + if (isStarted.getAndSet(true)) { + return; + } // This thread terminates on interruption or IO failure - messageProcessingThread = new Thread(new Runnable() { + Thread messageProcessingThread = new Thread(new Runnable() { public void run() { try { while (true) { @@ -274,7 +318,7 @@ messageProcessingThread.start(); // This thread only terminates if it is interrupted - sendThread = new Thread(new Runnable() { + Thread sendThread = new Thread(new Runnable() { public void run() { while (true) { try { @@ -293,33 +337,6 @@ sendThread.setDaemon(true); sendThread.start(); } - - /** - * Asynchronously executes the request on a remote server. - * - * @param requestMessage The request to execute - * - * @return a {...@link Future} that can be used to access the server's response - */ - public Future<Response> executeRequestAsync(final Request requestMessage) { - Future<Response> responseFuture = serverRequestExecutor.submit(new Callable<Response>() { - public Response call() throws Exception { - Message.Builder messageBuilder = Message.newBuilder(); - int messageId = nextMessageId.getAndIncrement(); - messageBuilder.setMessageId(messageId); - messageBuilder.setMessageType(Message.MessageType.REQUEST); - messageBuilder.setRequest(requestMessage); - - Message message = messageBuilder.build(); - PendingRequest pendingRequest = new PendingRequest(message); - sendQueue.put(pendingRequest); - - return pendingRequest.waitForResponse(); - } - }); - - return responseFuture; - } private void processClientRequest(int messageId, Request request) throws InterruptedException { ======================================= --- /trunk/dev/core/src/com/google/gwt/dev/shell/remoteui/RemoteUI.java Wed Nov 11 18:23:59 2009 +++ /trunk/dev/core/src/com/google/gwt/dev/shell/remoteui/RemoteUI.java Thu Nov 12 13:26:49 2009 @@ -62,6 +62,7 @@ devModeRequestProcessor = new DevModeServiceRequestProcessor(this); transport = new MessageTransport(transportSocket.getInputStream(), transportSocket.getOutputStream(), devModeRequestProcessor, this); + transport.start(); } catch (UnknownHostException e) { throw new RuntimeException(e); } catch (IOException e) { @@ -123,13 +124,11 @@ } public void onTermination(Exception e) { - getTopLogger().log( - TreeLogger.INFO, - "Remote UI connection terminated due to exception: " - + e); + getTopLogger().log(TreeLogger.INFO, + "Remote UI connection terminated due to exception: " + e); getTopLogger().log(TreeLogger.INFO, "Shutting down development mode server."); - + try { // Close the transport socket transportSocket.close(); ======================================= --- /trunk/dev/core/test/com/google/gwt/dev/shell/remoteui/MessageTransportTest.java Wed Nov 11 20:59:51 2009 +++ /trunk/dev/core/test/com/google/gwt/dev/shell/remoteui/MessageTransportTest.java Thu Nov 12 13:26:49 2009 @@ -117,6 +117,7 @@ public void onTermination(Exception e) { } }); + messageTransport.start(); Message.Request.Builder requestMessageBuilder = Message.Request.newBuilder(); requestMessageBuilder.setServiceType(Message.Request.ServiceType.DEV_MODE); @@ -170,6 +171,7 @@ MessageTransport messageTransport = new MessageTransport( network.getClientSocket().getInputStream(), network.getClientSocket().getOutputStream(), requestProcessor, null); + messageTransport.start(); // Generate a new request DevModeRequest.Builder devModeRequestBuilder = DevModeRequest.newBuilder(); @@ -245,6 +247,7 @@ public void onTermination(Exception e) { } }); + messageTransport.start(); // Generate a new request Message.Request.Builder requestMessageBuilder = Message.Request.newBuilder(); @@ -322,6 +325,7 @@ public void onTermination(Exception e) { } }); + messageTransport.start(); Message.Request.Builder requestMessageBuilder = Message.Request.newBuilder(); requestMessageBuilder.setServiceType(Message.Request.ServiceType.DEV_MODE); @@ -387,12 +391,14 @@ }; // Start up the message transport on the server side - new MessageTransport(network.getClientSocket().getInputStream(), + MessageTransport messageTransport = new MessageTransport( + network.getClientSocket().getInputStream(), network.getClientSocket().getOutputStream(), requestProcessor, new MessageTransport.TerminationCallback() { public void onTermination(Exception e) { } }); + messageTransport.start(); // Send the request from the client to the server Message.Builder clientRequestMsgBuilder = Message.newBuilder(); @@ -441,12 +447,14 @@ }; // Start up the message transport on the server side - new MessageTransport(network.getClientSocket().getInputStream(), + MessageTransport messageTransport = new MessageTransport( + network.getClientSocket().getInputStream(), network.getClientSocket().getOutputStream(), requestProcessor, new MessageTransport.TerminationCallback() { public void onTermination(Exception e) { } }); + messageTransport.start(); // Send a request to the server Message.Request.Builder clientRequestBuilder = Message.Request.newBuilder(); -- http://groups.google.com/group/Google-Web-Toolkit-Contributors
