Revision: 6881
Author: [email protected]
Date: Thu Nov 12 13:31:08 2009
Log: Merge tr...@r6880 into releases/2.0. Merge performed with the  
following command:

svn merge --ignore-ancestry -c 6880  
http://google-web-toolkit.googlecode.com/svn/trunk .

http://code.google.com/p/google-web-toolkit/source/detail?r=6881

Modified:
   
/releases/2.0/dev/core/src/com/google/gwt/dev/shell/remoteui/MessageTransport.java
  /releases/2.0/dev/core/src/com/google/gwt/dev/shell/remoteui/RemoteUI.java
   
/releases/2.0/dev/core/test/com/google/gwt/dev/shell/remoteui/MessageTransportTest.java

=======================================
---  
/releases/2.0/dev/core/src/com/google/gwt/dev/shell/remoteui/MessageTransport.java
       
Wed Nov 11 18:35:07 2009
+++  
/releases/2.0/dev/core/src/com/google/gwt/dev/shell/remoteui/MessageTransport.java
       
Thu Nov 12 13:31:08 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 {
=======================================
---  
/releases/2.0/dev/core/src/com/google/gwt/dev/shell/remoteui/RemoteUI.java      
 
Wed Nov 11 18:35:07 2009
+++  
/releases/2.0/dev/core/src/com/google/gwt/dev/shell/remoteui/RemoteUI.java      
 
Thu Nov 12 13:31:08 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();
=======================================
---  
/releases/2.0/dev/core/test/com/google/gwt/dev/shell/remoteui/MessageTransportTest.java
  
Wed Nov 11 21:11:00 2009
+++  
/releases/2.0/dev/core/test/com/google/gwt/dev/shell/remoteui/MessageTransportTest.java
  
Thu Nov 12 13:31:08 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

Reply via email to