[infinispan-commits] Infinispan SVN: r292 - in trunk/core/src/main: java/org/infinispan/config/parsing and 6 other directories.

infinispan-commits at lists.jboss.org infinispan-commits at lists.jboss.org
Thu May 14 06:55:29 EDT 2009


Author: manik.surtani at jboss.com
Date: 2009-05-14 06:55:29 -0400 (Thu, 14 May 2009)
New Revision: 292

Modified:
   trunk/core/src/main/java/org/infinispan/config/Configuration.java
   trunk/core/src/main/java/org/infinispan/config/parsing/XmlConfigurationParserImpl.java
   trunk/core/src/main/java/org/infinispan/factories/KnownComponentNames.java
   trunk/core/src/main/java/org/infinispan/factories/NamedExecutorsFactory.java
   trunk/core/src/main/java/org/infinispan/interceptors/base/BaseRpcInterceptor.java
   trunk/core/src/main/java/org/infinispan/remoting/rpc/CacheRpcManager.java
   trunk/core/src/main/java/org/infinispan/remoting/rpc/ResponseMode.java
   trunk/core/src/main/java/org/infinispan/remoting/rpc/RpcManagerImpl.java
   trunk/core/src/main/java/org/infinispan/remoting/transport/Transport.java
   trunk/core/src/main/java/org/infinispan/remoting/transport/jgroups/CommandAwareRpcDispatcher.java
   trunk/core/src/main/java/org/infinispan/remoting/transport/jgroups/JGroupsTransport.java
   trunk/core/src/main/resources/schema/infinispan-config-4.0.xsd
Log:
Added separate configuration for asynchronous marshalling

Modified: trunk/core/src/main/java/org/infinispan/config/Configuration.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/config/Configuration.java	2009-05-14 10:49:52 UTC (rev 291)
+++ trunk/core/src/main/java/org/infinispan/config/Configuration.java	2009-05-14 10:55:29 UTC (rev 292)
@@ -213,6 +213,7 @@
    private long rehashWaitTime = 60000;
    private boolean useLockStriping = true;
    private boolean unsafeUnreliableReturnValues = false;
+   private boolean asyncMarshalling = true;
 
    @Start(priority = 1)
    private void correctIsolationLevels() {
@@ -463,10 +464,18 @@
       this.rehashWaitTime = rehashWaitTime;
    }
 
+   public void setAsyncMarshalling(boolean asyncMarshalling) {
+      testImmutability("asyncMarshalling");
+      this.asyncMarshalling = asyncMarshalling;
+   }
+
    // ------------------------------------------------------------------------------------------------------------
    //   GETTERS
    // ------------------------------------------------------------------------------------------------------------
 
+   public boolean isAsyncMarshalling() {
+      return asyncMarshalling;
+   }
 
    public boolean isUseReplQueue() {
       return useReplQueue;
@@ -613,6 +622,7 @@
       if (isolationLevel != that.isolationLevel) return false;
       if (transactionManagerLookupClass != null ? !transactionManagerLookupClass.equals(that.transactionManagerLookupClass) : that.transactionManagerLookupClass != null)
          return false;
+      if (asyncMarshalling != that.asyncMarshalling) return false;
 
       return true;
    }
@@ -652,6 +662,7 @@
       result = 31 * result + (l1OnRehash ? 1 : 0);
       result = 31 * result + (consistentHashClass != null ? consistentHashClass.hashCode() : 0);
       result = 31 * result + numOwners;
+      result = 31 * result + (asyncMarshalling ? 1 : 0);
       return result;
    }
 

Modified: trunk/core/src/main/java/org/infinispan/config/parsing/XmlConfigurationParserImpl.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/config/parsing/XmlConfigurationParserImpl.java	2009-05-14 10:49:52 UTC (rev 291)
+++ trunk/core/src/main/java/org/infinispan/config/parsing/XmlConfigurationParserImpl.java	2009-05-14 10:55:29 UTC (rev 292)
@@ -347,6 +347,8 @@
       if (existsAttribute(tmp)) config.setReplQueueMaxElements(getInt(tmp));
       tmp = getAttributeValue(element, "useAsyncSerialization");
       if (existsAttribute(tmp)) config.setUseAsyncSerialization(getBoolean(tmp));
+      tmp = getAttributeValue(element, "asyncMarshalling");
+      if (existsAttribute(tmp)) config.setAsyncMarshalling(getBoolean(tmp));
    }
 
    void configureLocking(Element element, Configuration config) {

Modified: trunk/core/src/main/java/org/infinispan/factories/KnownComponentNames.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/factories/KnownComponentNames.java	2009-05-14 10:49:52 UTC (rev 291)
+++ trunk/core/src/main/java/org/infinispan/factories/KnownComponentNames.java	2009-05-14 10:55:29 UTC (rev 292)
@@ -8,7 +8,7 @@
  * @since 4.0
  */
 public class KnownComponentNames {
-   public static final String ASYNC_SERIALIZATION_EXECUTOR = "org.infinispan.executors.serialization";
+   public static final String ASYNC_TRANSPORT_EXECUTOR = "org.infinispan.executors.transport";
    public static final String ASYNC_NOTIFICATION_EXECUTOR = "org.infinispan.executors.notification";
    public static final String EVICTION_SCHEDULED_EXECUTOR = "org.infinispan.executors.eviction";
    public static final String ASYNC_REPLICATION_QUEUE_EXECUTOR = "org.infinispan.executors.replicationQueue";

Modified: trunk/core/src/main/java/org/infinispan/factories/NamedExecutorsFactory.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/factories/NamedExecutorsFactory.java	2009-05-14 10:49:52 UTC (rev 291)
+++ trunk/core/src/main/java/org/infinispan/factories/NamedExecutorsFactory.java	2009-05-14 10:55:29 UTC (rev 292)
@@ -26,7 +26,7 @@
          if (componentName.equals(KnownComponentNames.ASYNC_NOTIFICATION_EXECUTOR)) {
             return (T) buildAndConfigureExecutorService(globalConfiguration.getAsyncListenerExecutorFactoryClass(),
                                                         globalConfiguration.getAsyncListenerExecutorProperties());
-         } else if (componentName.equals(KnownComponentNames.ASYNC_SERIALIZATION_EXECUTOR)) {
+         } else if (componentName.equals(KnownComponentNames.ASYNC_TRANSPORT_EXECUTOR)) {
             return (T) buildAndConfigureExecutorService(globalConfiguration.getAsyncSerializationExecutorFactoryClass(),
                                                         globalConfiguration.getAsyncSerializationExecutorProperties());
          } else if (componentName.equals(KnownComponentNames.EVICTION_SCHEDULED_EXECUTOR)) {

Modified: trunk/core/src/main/java/org/infinispan/interceptors/base/BaseRpcInterceptor.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/interceptors/base/BaseRpcInterceptor.java	2009-05-14 10:49:52 UTC (rev 291)
+++ trunk/core/src/main/java/org/infinispan/interceptors/base/BaseRpcInterceptor.java	2009-05-14 10:55:29 UTC (rev 292)
@@ -50,7 +50,7 @@
 
    @Inject
    public void init(CacheRpcManager rpcManager,
-                    @ComponentName(KnownComponentNames.ASYNC_SERIALIZATION_EXECUTOR) ExecutorService e) {
+                    @ComponentName(KnownComponentNames.ASYNC_TRANSPORT_EXECUTOR) ExecutorService e) {
       this.rpcManager = rpcManager;
       this.asyncExecutorService = e;
    }

Modified: trunk/core/src/main/java/org/infinispan/remoting/rpc/CacheRpcManager.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/remoting/rpc/CacheRpcManager.java	2009-05-14 10:49:52 UTC (rev 291)
+++ trunk/core/src/main/java/org/infinispan/remoting/rpc/CacheRpcManager.java	2009-05-14 10:55:29 UTC (rev 292)
@@ -69,6 +69,10 @@
       }
    }
 
+   private ResponseMode getResponseMode(boolean sync) {
+      return sync ? ResponseMode.SYNCHRONOUS : configuration.isAsyncMarshalling() ? ResponseMode.ASYNCHRONOUS : ResponseMode.ASYNCHRONOUS_WITH_SYNC_MARSHALLING;
+   }
+
    public void multicastRpcCommand(List<Address> recipients, CacheRpcCommand command, boolean sync, boolean useOutOfBandMessage) throws ReplicationException {
       if (trace) {
          log.trace("invoking method " + command.getClass().getSimpleName() + ", members=" + rpcManager.getTransport().getMembers() + ", mode=" +
@@ -81,7 +85,7 @@
       try {
          rsps = rpcManager.invokeRemotely(recipients,
                                           command,
-                                          sync ? ResponseMode.SYNCHRONOUS : ResponseMode.ASYNCHRONOUS, // is synchronised?
+                                          getResponseMode(sync),
                                           configuration.getSyncReplTimeout(), useOutOfBandMessage, stateTransferEnabled
          );
          if (trace) log.trace("responses=" + rsps);

Modified: trunk/core/src/main/java/org/infinispan/remoting/rpc/ResponseMode.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/remoting/rpc/ResponseMode.java	2009-05-14 10:49:52 UTC (rev 291)
+++ trunk/core/src/main/java/org/infinispan/remoting/rpc/ResponseMode.java	2009-05-14 10:55:29 UTC (rev 292)
@@ -7,5 +7,12 @@
  * @since 4.0
  */
 public enum ResponseMode {
-   SYNCHRONOUS, ASYNCHRONOUS, WAIT_FOR_VALID_RESPONSE
+   SYNCHRONOUS,
+   ASYNCHRONOUS,
+   ASYNCHRONOUS_WITH_SYNC_MARSHALLING,
+   WAIT_FOR_VALID_RESPONSE;
+
+   public boolean isAsynchronous() {
+      return this == ASYNCHRONOUS || this == ASYNCHRONOUS_WITH_SYNC_MARSHALLING;
+   }
 }

Modified: trunk/core/src/main/java/org/infinispan/remoting/rpc/RpcManagerImpl.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/remoting/rpc/RpcManagerImpl.java	2009-05-14 10:49:52 UTC (rev 291)
+++ trunk/core/src/main/java/org/infinispan/remoting/rpc/RpcManagerImpl.java	2009-05-14 10:55:29 UTC (rev 292)
@@ -49,7 +49,7 @@
    @Inject
    public void injectDependencies(GlobalConfiguration globalConfiguration, Transport t, InboundInvocationHandler handler,
                                   Marshaller marshaller,
-                                  @ComponentName(KnownComponentNames.ASYNC_SERIALIZATION_EXECUTOR) ExecutorService e,
+                                  @ComponentName(KnownComponentNames.ASYNC_TRANSPORT_EXECUTOR) ExecutorService e,
                                   CacheManagerNotifier notifier) {
       this.t = t;
       this.t.initialize(globalConfiguration, globalConfiguration.getTransportProperties(), marshaller, e, handler,
@@ -71,7 +71,7 @@
          List<Response> result = t.invokeRemotely(recipients, rpcCommand, mode, timeout, usePriorityQueue, responseFilter, stateTransferEnabled);
          if (isStatisticsEnabled()) replicationCount.incrementAndGet();
          return result;
-      } catch (CacheException  e) {
+      } catch (CacheException e) {
          if (log.isTraceEnabled()) log.trace("replicaiton exception: ", e);
          if (isStatisticsEnabled()) replicationFailures.incrementAndGet();
          throw e;

Modified: trunk/core/src/main/java/org/infinispan/remoting/transport/Transport.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/remoting/transport/Transport.java	2009-05-14 10:49:52 UTC (rev 291)
+++ trunk/core/src/main/java/org/infinispan/remoting/transport/Transport.java	2009-05-14 10:55:29 UTC (rev 292)
@@ -9,9 +9,9 @@
 import org.infinispan.marshall.Marshaller;
 import org.infinispan.notifications.cachemanagerlistener.CacheManagerNotifier;
 import org.infinispan.remoting.InboundInvocationHandler;
+import org.infinispan.remoting.responses.Response;
 import org.infinispan.remoting.rpc.ResponseFilter;
 import org.infinispan.remoting.rpc.ResponseMode;
-import org.infinispan.remoting.responses.Response;
 import org.infinispan.statetransfer.StateTransferException;
 
 import java.util.List;
@@ -59,7 +59,8 @@
     * @return a list of responses from each member contacted.
     * @throws Exception in the event of problems.
     */
-   List<Response> invokeRemotely(List<Address> recipients, ReplicableCommand rpcCommand, ResponseMode mode, long timeout, boolean usePriorityQueue, ResponseFilter responseFilter, boolean supportReplay) throws Exception;
+   List<Response> invokeRemotely(List<Address> recipients, ReplicableCommand rpcCommand, ResponseMode mode, long timeout,
+                                 boolean usePriorityQueue, ResponseFilter responseFilter, boolean supportReplay) throws Exception;
 
    /**
     * @return true if the current Channel is the coordinator of the cluster.

Modified: trunk/core/src/main/java/org/infinispan/remoting/transport/jgroups/CommandAwareRpcDispatcher.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/remoting/transport/jgroups/CommandAwareRpcDispatcher.java	2009-05-14 10:49:52 UTC (rev 291)
+++ trunk/core/src/main/java/org/infinispan/remoting/transport/jgroups/CommandAwareRpcDispatcher.java	2009-05-14 10:55:29 UTC (rev 292)
@@ -96,15 +96,16 @@
     * org.jgroups.blocks.RspFilter)} except that this version is aware of {@link ReplicableCommand} objects.
     */
    public RspList invokeRemoteCommands(Vector<Address> dests, ReplicableCommand command, int mode, long timeout,
-                                       boolean anycasting, boolean oob, RspFilter filter, boolean supportReplay)
+                                       boolean anycasting, boolean oob, RspFilter filter, boolean supportReplay, boolean asyncMarshalling)
          throws NotSerializableException, ExecutionException, InterruptedException {
+
       ReplicationTask task = new ReplicationTask(command, oob, dests, mode, timeout, anycasting, filter, supportReplay);
 
-      if (mode == GroupRequest.GET_NONE) {
+      if (asyncMarshalling) {
          asyncExecutor.submit(task);
          return null; // don't wait for a response!
       } else {
-         RspList response = null;
+         RspList response;
          try {
             response = task.call();
          } catch (Exception e) {

Modified: trunk/core/src/main/java/org/infinispan/remoting/transport/jgroups/JGroupsTransport.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/remoting/transport/jgroups/JGroupsTransport.java	2009-05-14 10:49:52 UTC (rev 291)
+++ trunk/core/src/main/java/org/infinispan/remoting/transport/jgroups/JGroupsTransport.java	2009-05-14 10:55:29 UTC (rev 292)
@@ -8,10 +8,10 @@
 import org.infinispan.notifications.cachemanagerlistener.CacheManagerNotifier;
 import org.infinispan.remoting.InboundInvocationHandler;
 import org.infinispan.remoting.ReplicationException;
+import org.infinispan.remoting.responses.ExceptionResponse;
+import org.infinispan.remoting.responses.Response;
 import org.infinispan.remoting.rpc.ResponseFilter;
 import org.infinispan.remoting.rpc.ResponseMode;
-import org.infinispan.remoting.responses.ExceptionResponse;
-import org.infinispan.remoting.responses.Response;
 import org.infinispan.remoting.transport.Address;
 import org.infinispan.remoting.transport.DistributedSync;
 import org.infinispan.remoting.transport.Transport;
@@ -289,13 +289,13 @@
       boolean unlock = true;
       // if there is a FLUSH in progress, block till it completes
       flushTracker.blockUntilReleased(distributedSyncTimeout, MILLISECONDS);
-
+      boolean asyncMarshalling = mode == ResponseMode.ASYNCHRONOUS;
       try {
          RspList rsps = dispatcher.invokeRemoteCommands(toJGroupsAddressVector(recipients), rpcCommand, toJGroupsMode(mode),
                                                         timeout, false, usePriorityQueue,
-                                                        toJGroupsFilter(responseFilter), supportReplay);
+                                                        toJGroupsFilter(responseFilter), supportReplay, asyncMarshalling);
 
-         if (mode == ResponseMode.ASYNCHRONOUS) return Collections.emptyList();// async case
+         if (mode.isAsynchronous()) return Collections.emptyList();// async case
 
          if (trace)
             log.trace("Cache [{0}]: responses for command {1}:\n{2}", getAddress(), rpcCommand.getClass().getSimpleName(), rsps);
@@ -341,6 +341,7 @@
    private int toJGroupsMode(ResponseMode mode) {
       switch (mode) {
          case ASYNCHRONOUS:
+         case ASYNCHRONOUS_WITH_SYNC_MARSHALLING:
             return GroupRequest.GET_NONE;
          case SYNCHRONOUS:
             return GroupRequest.GET_ALL;

Modified: trunk/core/src/main/resources/schema/infinispan-config-4.0.xsd
===================================================================
--- trunk/core/src/main/resources/schema/infinispan-config-4.0.xsd	2009-05-14 10:49:52 UTC (rev 291)
+++ trunk/core/src/main/resources/schema/infinispan-config-4.0.xsd	2009-05-14 10:55:29 UTC (rev 292)
@@ -19,7 +19,7 @@
          <xs:element name="globalJmxStatistics" type="tns:globalJmxStatisticsType" minOccurs="0" maxOccurs="1"/>
          <xs:element name="transport" type="tns:transportConfigType" minOccurs="0" maxOccurs="1"/>
          <xs:element name="asyncListenerExecutor" type="tns:executorConfigurationType" minOccurs="0" maxOccurs="1"/>
-         <xs:element name="asyncSerializationExecutor" type="tns:executorConfigurationType" minOccurs="0"
+         <xs:element name="asyncTransportExecutor" type="tns:executorConfigurationType" minOccurs="0"
                      maxOccurs="1"/>
          <xs:element name="evictionScheduledExecutor" type="tns:executorConfigurationType" minOccurs="0"
                      maxOccurs="1"/>
@@ -176,6 +176,7 @@
    </xs:complexType>
 
    <xs:complexType name="asyncType">
+      <xs:attribute name="asyncMarshalling" type="tns:booleanType"/>
       <xs:attribute name="useReplQueue" type="tns:booleanType"/>
       <xs:attribute name="replQueueInterval" type="tns:positiveNumber"/>
       <xs:attribute name="replQueueMaxElements" type="tns:positiveNumber"/>




More information about the infinispan-commits mailing list