[infinispan-commits] Infinispan SVN: r1708 - in trunk/server: core/src/main/scala/org/infinispan/server/core/transport/netty and 4 other directories.

infinispan-commits at lists.jboss.org infinispan-commits at lists.jboss.org
Wed Apr 21 10:58:10 EDT 2010


Author: galder.zamarreno at jboss.com
Date: 2010-04-21 10:58:08 -0400 (Wed, 21 Apr 2010)
New Revision: 1708

Added:
   trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/IdleStateHandlerProvider.scala
   trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/HotRodIdleTimeoutTest.scala
Modified:
   trunk/server/core/src/main/scala/org/infinispan/server/core/AbstractProtocolServer.scala
   trunk/server/core/src/main/scala/org/infinispan/server/core/Main.scala
   trunk/server/core/src/main/scala/org/infinispan/server/core/ProtocolServer.scala
   trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/DecoderAdapter.scala
   trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/EncoderAdapter.scala
   trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/NettyChannelPipelineFactory.scala
   trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/NettyTransport.scala
   trunk/server/hotrod/src/main/scala/org/infinispan/server/hotrod/HotRodServer.scala
   trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/HotRodConcurrentTest.scala
   trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/HotRodFunctionalTest.scala
   trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/HotRodReplicationTest.scala
   trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/HotRodSingleNodeTest.scala
   trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/test/HotRodClient.scala
   trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/test/HotRodTestingUtil.scala
   trunk/server/memcached/src/test/scala/org/infinispan/server/memcached/test/MemcachedTestingUtil.scala
Log:
[ISPN-385] (Add idle timeout to memcached and hot rod servers) Added -i command line parameter for controlling idle timeout.

Modified: trunk/server/core/src/main/scala/org/infinispan/server/core/AbstractProtocolServer.scala
===================================================================
--- trunk/server/core/src/main/scala/org/infinispan/server/core/AbstractProtocolServer.scala	2010-04-21 13:47:30 UTC (rev 1707)
+++ trunk/server/core/src/main/scala/org/infinispan/server/core/AbstractProtocolServer.scala	2010-04-21 14:58:08 UTC (rev 1708)
@@ -18,10 +18,8 @@
    private var workerThreads: Int = _
    private var transport: Transport = _
    private var cacheManager: CacheManager = _
-   private var decoder: Decoder = _
-   private var encoder: Encoder = _
 
-   override def start(host: String, port: Int, cacheManager: CacheManager, masterThreads: Int, workerThreads: Int) {
+   override def start(host: String, port: Int, cacheManager: CacheManager, masterThreads: Int, workerThreads: Int, idleTimeout: Int) {
       this.host = host
       this.port = port
       this.masterThreads = masterThreads
@@ -29,12 +27,10 @@
       this.cacheManager = cacheManager
 
       cacheManager.addListener(getRankCalculatorListener)
-      encoder = getEncoder
-      // TODO: add an IdleStateHandler so that idle connections are detected, this could help on malformed data
-      // TODO: ... requests such as when the lenght of data is bigger than the expected data itself.
-      val nettyEncoder = if (encoder != null) new EncoderAdapter(encoder) else null
       val address =  new InetSocketAddress(host, port)
-      transport = new NettyTransport(this, nettyEncoder, address, masterThreads, workerThreads, threadNamePrefix)
+      val encoder = getEncoder
+      val nettyEncoder = if (encoder != null) new EncoderAdapter(encoder) else null
+      transport = new NettyTransport(this, nettyEncoder, address, masterThreads, workerThreads, idleTimeout, threadNamePrefix)
       transport.start
    }
 
@@ -47,4 +43,4 @@
 
    def getPort = port
 
-}
\ No newline at end of file
+}

Modified: trunk/server/core/src/main/scala/org/infinispan/server/core/Main.scala
===================================================================
--- trunk/server/core/src/main/scala/org/infinispan/server/core/Main.scala	2010-04-21 13:47:30 UTC (rev 1707)
+++ trunk/server/core/src/main/scala/org/infinispan/server/core/Main.scala	2010-04-21 14:58:08 UTC (rev 1708)
@@ -24,10 +24,12 @@
    val PROP_KEY_WORKER_THREADS = "infinispan.server.worker.threads"
    val PROP_KEY_CACHE_CONFIG = "infinispan.server.cache.config"
    val PROP_KEY_PROTOCOL = "infinispan.server.protocol"
+   val PROP_KEY_IDLE_TIMEOUT = "infinispan.server.idle.timeout"
    val PORT_DEFAULT = 11211
    val HOST_DEFAULT = "127.0.0.1"
    val MASTER_THREADS_DEFAULT = 0
    val WORKER_THREADS_DEFAULT = 0
+   val IDLE_TIMEOUT_DEFAULT = 60
 
    /**
     * Server properties.  This object holds all of the required
@@ -36,7 +38,7 @@
     */
    private val props: mutable.Map[String, String] = {
       // Set default properties
-      val properties = new HashMap[String, String]()
+      val properties = new HashMap[String, String]
       val sysProps = System.getProperties
       for (property <- asIterator(sysProps.iterator))
          properties.put(property._1, property._2)
@@ -47,50 +49,53 @@
    private var server: ProtocolServer = _
 
    def main(args: Array[String]) {
-       info("Start main with args: {0}", args.mkString(", "))
-       var worker = new Callable[Void] {
-          override def call = {
-             try {
-                boot(args)
-             }
-             catch {
-                case e: Exception => {
-                   System.err.println("Failed to boot JBoss:")
-                   e.printStackTrace
-                   throw e
-                }
-             }
-             null
-          }
-       }
-       var f = Executors.newSingleThreadScheduledExecutor(new ThreadFactory {
-          override def newThread(r: Runnable): Thread = {
-             // TODO Maybe create thread names based on the protocol run
-             return new Thread(r, "InfinispanServer-Main")
-          }
-       }).submit(worker)
-       f.get
-    }
+      info("Start main with args: {0}", args.mkString(", "))
+      var worker = new Callable[Void] {
+         override def call = {
+            try {
+               boot(args)
+            }
+            catch {
+               case e: Exception => {
+                  System.err.println("Failed to boot JBoss:")
+                  e.printStackTrace
+                  throw e
+               }
+            }
+            null
+         }
+      }
+      var f = Executors.newSingleThreadScheduledExecutor(new ThreadFactory {
+         override def newThread(r: Runnable): Thread = {
+            // TODO Maybe create thread names based on the protocol run
+            return new Thread(r, "InfinispanServer-Main")
+         }
+      }).submit(worker)
+      f.get
+   }
 
-    def boot(args: Array[String]) {
+   def boot(args: Array[String]) {
       // First process the command line to pickup custom props/settings
       processCommandLine(args)
 
       val host = if (props.get(PROP_KEY_HOST) == None) HOST_DEFAULT else props.get(PROP_KEY_HOST).get
       val masterThreads = if (props.get(PROP_KEY_MASTER_THREADS) == None) MASTER_THREADS_DEFAULT else props.get(PROP_KEY_MASTER_THREADS).get.toInt
-      if (masterThreads < 0) {
+      if (masterThreads < 0)
          throw new IllegalArgumentException("Master threads can't be lower than 0: " + masterThreads)
-      }
+
       val workerThreads = if (props.get(PROP_KEY_WORKER_THREADS) == None) WORKER_THREADS_DEFAULT else props.get(PROP_KEY_WORKER_THREADS).get.toInt
-      if (workerThreads < 0) {
+      if (workerThreads < 0)
          throw new IllegalArgumentException("Worker threads can't be lower than 0: " + masterThreads)
-      }
+
       val configFile = props.get(PROP_KEY_CACHE_CONFIG)
       var protocol = props.get(PROP_KEY_PROTOCOL)
       if (protocol == None) {
          System.err.println("ERROR: Please indicate protocol to run with -r parameter")
          showAndExit
       }
+      val idleTimeout = if (props.get(PROP_KEY_IDLE_TIMEOUT) == None) IDLE_TIMEOUT_DEFAULT else props.get(PROP_KEY_IDLE_TIMEOUT).get.toInt
+      if (idleTimeout < 0)
+         throw new IllegalArgumentException("Idle timeout can't be lower than 0: " + idleTimeout)
 
       // TODO: move class name and protocol number to a resource file under the corresponding project
       val clazz = protocol.get match {
@@ -111,7 +116,7 @@
       var server = Util.getInstance(clazz).asInstanceOf[ProtocolServer]
       val cacheManager = if (configFile == None) new DefaultCacheManager else new DefaultCacheManager(configFile.get)
       addShutdownHook(new ShutdownHook(server, cacheManager))
-      server.start(host, port, cacheManager, masterThreads, workerThreads)
+      server.start(host, port, cacheManager, masterThreads, workerThreads, idleTimeout)
    }
 
    private def processCommandLine(args: Array[String]) {
@@ -125,7 +130,8 @@
          new LongOpt("master_threads", LongOpt.REQUIRED_ARGUMENT, null, 'm'),
          new LongOpt("worker_threads", LongOpt.REQUIRED_ARGUMENT, null, 't'),
          new LongOpt("cache_config", LongOpt.REQUIRED_ARGUMENT, null, 'c'),
-         new LongOpt("protocol", LongOpt.REQUIRED_ARGUMENT, null, 'r'))
+         new LongOpt("protocol", LongOpt.REQUIRED_ARGUMENT, null, 'r'),
+         new LongOpt("idle_timeout", LongOpt.REQUIRED_ARGUMENT, null, 'i'))
       var getopt = new Getopt(programName, args, sopts, lopts)
       var code: Int = 0
       while ((({code = getopt.getopt; code})) != -1) {
@@ -143,6 +149,7 @@
             case 't' => props.put(PROP_KEY_WORKER_THREADS, getopt.getOptarg)
             case 'c' => props.put(PROP_KEY_CACHE_CONFIG, getopt.getOptarg)
             case 'r' => props.put(PROP_KEY_PROTOCOL, getopt.getOptarg)
+            case 'i' => props.put(PROP_KEY_IDLE_TIMEOUT, getopt.getOptarg)
             case 'D' => {
                val arg = getopt.getOptarg
                var name = ""
@@ -184,6 +191,7 @@
       println("    -t, --work_threads=<num>           Number of threads processing incoming requests and sending responses (default: unlimited while resources are available)")
       println("    -c, --cache_config=<filename>      Cache configuration file (default: creates cache with default values)")
       println("    -r, --protocol=[memcached|hotrod]  Protocol to understand by the server. This is a mandatory option and you should choose one of the two options")
+      println("    -i, --idle_timeout=<num>           Idle read timeout used to detect stale connections (default: 60 seconds). If no new messages have been read within this time, the server disconnects the channel. Passing 0 means disabling idle timeout.")
       println("    -D<name>[=<value>]                 Set a system property")
       println
       System.exit(0)
@@ -210,4 +218,4 @@
          }
       }
    }
-}
\ No newline at end of file
+}

Modified: trunk/server/core/src/main/scala/org/infinispan/server/core/ProtocolServer.scala
===================================================================
--- trunk/server/core/src/main/scala/org/infinispan/server/core/ProtocolServer.scala	2010-04-21 13:47:30 UTC (rev 1707)
+++ trunk/server/core/src/main/scala/org/infinispan/server/core/ProtocolServer.scala	2010-04-21 14:58:08 UTC (rev 1708)
@@ -9,8 +9,8 @@
  * @since 4.1
  */
 trait ProtocolServer {
-   def start(host: String, port: Int, cacheManager: CacheManager, masterThreads: Int, workerThreads: Int)
+   def start(host: String, port: Int, cacheManager: CacheManager, masterThreads: Int, workerThreads: Int, idleTimeout: Int)
    def stop
    def getEncoder: Encoder
    def getDecoder: Decoder
-}
\ No newline at end of file
+}

Modified: trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/DecoderAdapter.scala
===================================================================
--- trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/DecoderAdapter.scala	2010-04-21 13:47:30 UTC (rev 1707)
+++ trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/DecoderAdapter.scala	2010-04-21 14:58:08 UTC (rev 1708)
@@ -28,6 +28,7 @@
 
    override def channelOpen(ctx: NettyChannelHandlerContext, e: ChannelStateEvent) {
       transport.acceptedChannels.add(e.getChannel)
+      super.channelOpen(ctx, e)
    }
 
 }
\ No newline at end of file

Modified: trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/EncoderAdapter.scala
===================================================================
--- trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/EncoderAdapter.scala	2010-04-21 13:47:30 UTC (rev 1707)
+++ trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/EncoderAdapter.scala	2010-04-21 14:58:08 UTC (rev 1708)
@@ -11,7 +11,6 @@
  * @author Galder Zamarreño
  * @since 4.1
  */
- at ChannelHandler.Sharable
 class EncoderAdapter(encoder: Encoder) extends OneToOneEncoder {
 
    protected override def encode(nCtx: NettyChannelHandlerContext, ch: NettyChannel, msg: AnyRef): AnyRef = {

Added: trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/IdleStateHandlerProvider.scala
===================================================================
--- trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/IdleStateHandlerProvider.scala	                        (rev 0)
+++ trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/IdleStateHandlerProvider.scala	2010-04-21 14:58:08 UTC (rev 1708)
@@ -0,0 +1,17 @@
+package org.infinispan.server.core.transport
+
+import org.jboss.netty.handler.timeout.{IdleStateEvent, IdleStateAwareChannelHandler}
+import org.jboss.netty.channel.{ChannelHandlerContext => NettyChannelHandlerContext}
+/**
+ * // TODO: Document this
+ * @author Galder Zamarreño
+ * @since // TODO
+ */
+class IdleStateHandlerProvider extends IdleStateAwareChannelHandler {
+
+   override def channelIdle(nCtx: NettyChannelHandlerContext, e: IdleStateEvent) {
+      nCtx.getChannel.disconnect
+      super.channelIdle(nCtx, e)
+   }
+
+}
\ No newline at end of file

Modified: trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/NettyChannelPipelineFactory.scala
===================================================================
--- trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/NettyChannelPipelineFactory.scala	2010-04-21 13:47:30 UTC (rev 1707)
+++ trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/NettyChannelPipelineFactory.scala	2010-04-21 14:58:08 UTC (rev 1708)
@@ -2,21 +2,35 @@
 
 import org.jboss.netty.channel._
 import org.infinispan.server.core.ProtocolServer
+import org.jboss.netty.handler.timeout.IdleStateHandler
+import org.infinispan.server.core.transport.IdleStateHandlerProvider
+import org.jboss.netty.util.{HashedWheelTimer, Timer}
 
 /**
  * // TODO: Document this
  * @author Galder Zamarreño
  * @since 4.1
  */
-class NettyChannelPipelineFactory(server: ProtocolServer, encoder: ChannelDownstreamHandler, transport: NettyTransport)
+class NettyChannelPipelineFactory(server: ProtocolServer, encoder: ChannelDownstreamHandler,
+                                  transport: NettyTransport, idleTimeout: Int)
       extends ChannelPipelineFactory {
 
+   private var timer: Timer = _
+
    override def getPipeline: ChannelPipeline = {
       val pipeline = Channels.pipeline
       pipeline.addLast("decoder", new DecoderAdapter(server.getDecoder, transport))
       if (encoder != null)
          pipeline.addLast("encoder", encoder)
+      if (idleTimeout != 0) {
+         timer = new HashedWheelTimer
+         pipeline.addLast("idleHandler", new IdleStateHandler(timer, idleTimeout, 0, 0))
+         pipeline.addLast("idleHandlerProvider", new IdleStateHandlerProvider)
+      }
       return pipeline;
    }
 
-}
\ No newline at end of file
+   def stop {
+      if (timer != null) timer.stop
+   }
+}

Modified: trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/NettyTransport.scala
===================================================================
--- trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/NettyTransport.scala	2010-04-21 13:47:30 UTC (rev 1707)
+++ trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/NettyTransport.scala	2010-04-21 14:58:08 UTC (rev 1708)
@@ -10,7 +10,7 @@
 import org.infinispan.server.core.transport.Transport
 import scala.collection.JavaConversions._
 import org.infinispan.server.core.{ProtocolServer, Logging}
-import org.jboss.netty.util.{ThreadNameDeterminer, ThreadRenamingRunnable}
+import org.jboss.netty.util.{HashedWheelTimer, ThreadNameDeterminer, ThreadRenamingRunnable}
 
 /**
  * // TODO: Document this
@@ -18,13 +18,13 @@
  * @since 4.1
  */
 class NettyTransport(server: ProtocolServer, encoder: ChannelDownstreamHandler,
-                  address: SocketAddress, masterThreads: Int, workerThreads: Int,
-                  threadNamePrefix: String) extends Transport {
+                     address: SocketAddress, masterThreads: Int, workerThreads: Int,
+                     idleTimeout: Int, threadNamePrefix: String) extends Transport {
    import NettyTransport._
 
    private val serverChannels = new DefaultChannelGroup(threadNamePrefix + "-Channels")
    val acceptedChannels = new DefaultChannelGroup(threadNamePrefix + "-Accepted")
-   private val pipeline = new NettyChannelPipelineFactory(server, encoder, this)
+   private val pipeline = new NettyChannelPipelineFactory(server, encoder, this, idleTimeout)
    private val factory = {
       if (workerThreads == 0)
          new NioServerSocketChannelFactory(masterExecutor, workerExecutor)
@@ -90,6 +90,7 @@
             }
          }
       }
+      pipeline.stop
       debug("Channel group completely closed, release external resources");
       factory.releaseExternalResources();
    }

Modified: trunk/server/hotrod/src/main/scala/org/infinispan/server/hotrod/HotRodServer.scala
===================================================================
--- trunk/server/hotrod/src/main/scala/org/infinispan/server/hotrod/HotRodServer.scala	2010-04-21 13:47:30 UTC (rev 1707)
+++ trunk/server/hotrod/src/main/scala/org/infinispan/server/hotrod/HotRodServer.scala	2010-04-21 14:58:08 UTC (rev 1708)
@@ -22,8 +22,8 @@
 
    override def getDecoder: Decoder = new HotRodDecoder(getCacheManager)
 
-   override def start(host: String, port: Int, cacheManager: CacheManager, masterThreads: Int, workerThreads: Int) {
-      super.start(host, port, cacheManager, masterThreads, workerThreads)
+   override def start(host: String, port: Int, cacheManager: CacheManager, masterThreads: Int, workerThreads: Int, idleTimeout: Int) {
+      super.start(host, port, cacheManager, masterThreads, workerThreads, idleTimeout)
       // If clustered, set up a cache for topology information
       if (cacheManager.getGlobalConfiguration.getTransportClass != null) {
          defineTopologyCacheConfig(cacheManager)

Modified: trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/HotRodConcurrentTest.scala
===================================================================
--- trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/HotRodConcurrentTest.scala	2010-04-21 13:47:30 UTC (rev 1707)
+++ trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/HotRodConcurrentTest.scala	2010-04-21 14:58:08 UTC (rev 1708)
@@ -4,8 +4,6 @@
 import java.util.concurrent.{Callable, Executors, Future, CyclicBarrier}
 import test.HotRodClient
 import org.testng.annotations.Test
-import org.infinispan.test.fwk.TestCacheManagerFactory
-import org.infinispan.manager.CacheManager
 
 /**
  * // TODO: Document this
@@ -42,9 +40,9 @@
 
    class Operator(barrier: CyclicBarrier, m: Method, clientId: Int, numOpsPerClient: Int) extends Callable[Unit] {
 
-      private lazy val client = new HotRodClient("127.0.0.1", server.getPort, cacheName)
+      private lazy val client = new HotRodClient("127.0.0.1", server.getPort, cacheName, 60)
 
-      def call {
+      override def call {
          log.debug("Wait for all executions paths to be ready to perform calls", null)
          barrier.await
          try {

Modified: trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/HotRodFunctionalTest.scala
===================================================================
--- trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/HotRodFunctionalTest.scala	2010-04-21 13:47:30 UTC (rev 1707)
+++ trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/HotRodFunctionalTest.scala	2010-04-21 14:58:08 UTC (rev 1708)
@@ -1,12 +1,11 @@
 package org.infinispan.server.hotrod
 
-import org.infinispan.test.fwk.TestCacheManagerFactory
 import org.testng.annotations.Test
 import java.lang.reflect.Method
 import test.HotRodTestingUtil._
 import org.testng.Assert._
 import java.util.Arrays
-import org.infinispan.manager.{DefaultCacheManager, CacheManager}
+import org.infinispan.manager.DefaultCacheManager
 import org.infinispan.server.core.CacheValue
 import org.infinispan.server.hotrod.OperationStatus._
 

Added: trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/HotRodIdleTimeoutTest.scala
===================================================================
--- trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/HotRodIdleTimeoutTest.scala	                        (rev 0)
+++ trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/HotRodIdleTimeoutTest.scala	2010-04-21 14:58:08 UTC (rev 1708)
@@ -0,0 +1,38 @@
+package org.infinispan.server.hotrod
+
+import org.testng.annotations.Test
+import java.lang.reflect.Method
+import test.HotRodTestingUtil._
+import org.testng.Assert._
+import org.infinispan.manager.CacheManager
+import test.{HotRodClient, UniquePortThreadLocal}
+
+/**
+ * // TODO: Document this
+ * @author Galder Zamarreño
+ * @since // TODO
+ */
+ at Test(groups = Array("functional"), testName = "server.hotrod.HotRodIdleTimeoutTest")
+class HotRodIdleTimeoutTest extends HotRodSingleNodeTest {
+
+   override protected def createStartHotRodServer(cacheManager: CacheManager) =
+      startHotRodServer(cacheManager, UniquePortThreadLocal.get.intValue, 5)
+
+   override protected def connectClient = new HotRodClient("127.0.0.1", server.getPort, cacheName, 10)
+
+   def testSendPartialRequest(m: Method) {
+      client.assertPut(m)
+      val resp = client.executePartial(0xA0, 0x03, cacheName, k(m) , 0, 0, v(m), 0)
+      assertNull(resp) // No response received within expected timeout.
+      client.assertPutFail(m)
+      shutdownClient
+      
+      val newClient = connectClient
+      try {
+         newClient.assertPut(m)
+      } finally {
+         shutdownClient
+      }
+   }
+
+}
\ No newline at end of file

Modified: trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/HotRodReplicationTest.scala
===================================================================
--- trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/HotRodReplicationTest.scala	2010-04-21 13:47:30 UTC (rev 1707)
+++ trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/HotRodReplicationTest.scala	2010-04-21 14:58:08 UTC (rev 1708)
@@ -16,7 +16,7 @@
  * @since
  */
 
- at Test(groups = Array("functional"), testName = "server.hotrod.ClusterTest")
+ at Test(groups = Array("functional"), testName = "server.hotrod.HotRodReplicationTest")
 class HotRodReplicationTest extends MultipleCacheManagersTest {
 
    import HotRodServer._
@@ -27,25 +27,15 @@
 
    @Test(enabled=false) // Disable explicitly to avoid TestNG thinking this is a test!!
    override def createCacheManagers {
-      val config = getDefaultClusteredConfig(Configuration.CacheMode.REPL_SYNC)
-      config.setFetchInMemoryState(true)
-
-      val topologyCacheConfig = new Configuration
-      topologyCacheConfig.setCacheMode(CacheMode.REPL_SYNC)
-      topologyCacheConfig.setSyncReplTimeout(10000) // Milliseconds
-      topologyCacheConfig.setFetchInMemoryState(true) // State transfer required
-      topologyCacheConfig.setSyncCommitPhase(true) // Only for testing, so that asserts work fine.
-      topologyCacheConfig.setSyncRollbackPhase(true) // Only for testing, so that asserts work fine.
-
       for (i <- 0 until 2) {
          val cm = addClusterEnabledCacheManager()
-         cm.defineConfiguration(cacheName, config)
-         cm.defineConfiguration(TopologyCacheName, topologyCacheConfig)
+         cm.defineConfiguration(cacheName, createCacheConfig)
+         cm.defineConfiguration(TopologyCacheName, createTopologyCacheConfig)
       }
-      servers = startHotRodServer(cacheManagers.get(0)) :: servers
-      servers = startHotRodServer(cacheManagers.get(1), servers.head.getPort + 50) :: servers
+      servers = servers ::: List(startHotRodServer(cacheManagers.get(0))) 
+      servers = servers ::: List(startHotRodServer(cacheManagers.get(1), servers.head.getPort + 50))
       servers.foreach {s =>
-         clients = new HotRodClient("127.0.0.1", s.getPort, cacheName) :: clients
+         clients = new HotRodClient("127.0.0.1", s.getPort, cacheName, 60) :: clients
       }
    }
 
@@ -62,6 +52,22 @@
       // Do not clear cache between methods so that topology cache does not get cleared
    }
 
+   private def createCacheConfig: Configuration = {
+      val config = getDefaultClusteredConfig(Configuration.CacheMode.REPL_SYNC)
+      config.setFetchInMemoryState(true)
+      config
+   }
+
+   private def createTopologyCacheConfig: Configuration = {
+      val topologyCacheConfig = new Configuration
+      topologyCacheConfig.setCacheMode(CacheMode.REPL_SYNC)
+      topologyCacheConfig.setSyncReplTimeout(10000) // Milliseconds
+      topologyCacheConfig.setFetchInMemoryState(true) // State transfer required
+      topologyCacheConfig.setSyncCommitPhase(true) // Only for testing, so that asserts work fine.
+      topologyCacheConfig.setSyncRollbackPhase(true) // Only for testing, so that asserts work fine.
+      topologyCacheConfig
+   }
+
    def testReplicatedPut(m: Method) {
       val putSt = clients.head.put(k(m) , 0, 0, v(m)).status
       assertStatus(putSt, Success)
@@ -112,8 +118,8 @@
    private def assertTopologyReceived(topologyResp: AbstractTopologyResponse) {
       assertEquals(topologyResp.view.topologyId, 2)
       assertEquals(topologyResp.view.members.size, 2)
-      assertEquals(topologyResp.view.members.head, TopologyAddress("127.0.0.1", 11311, 0))
-      assertEquals(topologyResp.view.members.tail.head, TopologyAddress("127.0.0.1", 11361, 0))
+      assertEquals(topologyResp.view.members.head, TopologyAddress("127.0.0.1", servers.head.getPort, 0))
+      assertEquals(topologyResp.view.members.tail.head, TopologyAddress("127.0.0.1", servers.tail.head.getPort, 0))
    }
 
    def testReplicatedPutWithTopologyAwareClient(m: Method) {
@@ -131,6 +137,23 @@
       assertStatus(resp.status, Success)
       assertEquals(resp.topologyResponse, None)
       assertSuccess(clients.tail.head.get(k(m), 0), v(m, "v3-"))
+
+      val cm = addClusterEnabledCacheManager()
+      cm.defineConfiguration(cacheName, createCacheConfig)
+      cm.defineConfiguration(TopologyCacheName, createTopologyCacheConfig)
+      servers = servers ::: List(startHotRodServer(cacheManagers.get(2), servers.tail.head.getPort + 25)) 
+
+      resp = clients.head.put(k(m) , 0, 0, v(m, "v4-"), 2, 2)
+      assertStatus(resp.status, Success)
+      assertEquals(resp.topologyResponse.get.view.topologyId, 3)
+      assertEquals(resp.topologyResponse.get.view.members.size, 3)
+      assertEquals(resp.topologyResponse.get.view.members.head, TopologyAddress("127.0.0.1", servers.head.getPort, 0))
+      assertEquals(resp.topologyResponse.get.view.members.tail.head, TopologyAddress("127.0.0.1", servers.tail.head.getPort, 0))
+      assertEquals(resp.topologyResponse.get.view.members.tail.tail.head, TopologyAddress("127.0.0.1", servers.tail.tail.head.getPort, 0))
+      assertSuccess(clients.tail.head.get(k(m), 0), v(m, "v4-"))
+
+//      // TODO: Add stopping a server
+//      servers.tail.tail.head.stop
    }
 
 }
\ No newline at end of file

Modified: trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/HotRodSingleNodeTest.scala
===================================================================
--- trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/HotRodSingleNodeTest.scala	2010-04-21 13:47:30 UTC (rev 1707)
+++ trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/HotRodSingleNodeTest.scala	2010-04-21 14:58:08 UTC (rev 1708)
@@ -25,13 +25,15 @@
    override def createCacheManager: CacheManager = {
       val cacheManager = createTestCacheManager
       advancedCache = cacheManager.getCache[CacheKey, CacheValue](cacheName).getAdvancedCache
-      hotRodServer = startHotRodServer(cacheManager)
-      hotRodClient = new HotRodClient("127.0.0.1", hotRodServer.getPort, cacheName)
+      hotRodServer = createStartHotRodServer(cacheManager)
+      hotRodClient = connectClient
       cacheManager
    }
 
-   protected def createTestCacheManager: CacheManager = TestCacheManagerFactory.createLocalCacheManager(true) 
+   protected def createTestCacheManager: CacheManager = TestCacheManagerFactory.createLocalCacheManager(true)
 
+   protected def createStartHotRodServer(cacheManager: CacheManager) = startHotRodServer(cacheManager)
+
    @AfterClass(alwaysRun = true)
    override def destroyAfterClass {
       log.debug("Test finished, close cache, client and Hot Rod server", null)
@@ -47,4 +49,6 @@
    protected def jmxDomain = hotRodJmxDomain
 
    protected def shutdownClient: ChannelFuture = hotRodClient.stop
+
+   protected def connectClient: HotRodClient = new HotRodClient("127.0.0.1", hotRodServer.getPort, cacheName, 60)
 }
\ No newline at end of file

Modified: trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/test/HotRodClient.scala
===================================================================
--- trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/test/HotRodClient.scala	2010-04-21 13:47:30 UTC (rev 1707)
+++ trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/test/HotRodClient.scala	2010-04-21 14:58:08 UTC (rev 1708)
@@ -33,13 +33,13 @@
  * @author Galder Zamarreño
  * @since 4.1
  */
-class HotRodClient(host: String, port: Int, defaultCacheName: String) {
+class HotRodClient(host: String, port: Int, defaultCacheName: String, rspTimeoutSeconds: Int) {
    val idToOp = new ConcurrentHashMap[Long, Op]    
 
    private lazy val ch: Channel = {
       val factory = new NioClientSocketChannelFactory(Executors.newCachedThreadPool, Executors.newCachedThreadPool)
       val bootstrap: ClientBootstrap = new ClientBootstrap(factory)
-      bootstrap.setPipelineFactory(new ClientPipelineFactory(this))
+      bootstrap.setPipelineFactory(new ClientPipelineFactory(this, rspTimeoutSeconds))
       bootstrap.setOption("tcpNoDelay", true)
       bootstrap.setOption("keepAlive", true)
       // Make a new connection.
@@ -63,6 +63,14 @@
       assertStatus(status, Success)
    }
 
+   def assertPutFail(m: Method) {
+      val op = new Op(0xA0, 0x01, defaultCacheName, k(m), 0, 0, v(m), 0, 1 , 0, 0)
+      idToOp.put(op.id, op)
+      val future = ch.write(op)
+      future.awaitUninterruptibly
+      assertFalse(future.isSuccess)
+   }
+
    def assertPut(m: Method, kPrefix: String, vPrefix: String) {
       val status = put(k(m, kPrefix) , 0, 0, v(m, vPrefix)).status
       assertStatus(status, Success)
@@ -118,6 +126,12 @@
       execute(op, 0).asInstanceOf[ErrorResponse]
    }
 
+   def executePartial(magic: Int, code: Byte, name: String, k: Array[Byte], lifespan: Int, maxIdle: Int,
+                      v: Array[Byte], version: Long): ErrorResponse = {
+      val op = new PartialOp(magic, code, name, k, lifespan, maxIdle, v, 0, version, 1, 0)
+      execute(op, op.id).asInstanceOf[ErrorResponse]
+   }
+
    def execute(magic: Int, code: Byte, name: String, k: Array[Byte], lifespan: Int, maxIdle: Int,
                v: Array[Byte], version: Long, flags: Int): Response = {
       val op = new Op(magic, code, name, k, lifespan, maxIdle, v, flags, version, 1, 0)
@@ -182,13 +196,13 @@
 
 }
 
-private class ClientPipelineFactory(client: HotRodClient) extends ChannelPipelineFactory {
+private class ClientPipelineFactory(client: HotRodClient, rspTimeoutSeconds: Int) extends ChannelPipelineFactory {
 
    override def getPipeline = {
       val pipeline = Channels.pipeline
       pipeline.addLast("decoder", new Decoder(client))
       pipeline.addLast("encoder", new Encoder)
-      pipeline.addLast("handler", new ClientHandler)
+      pipeline.addLast("handler", new ClientHandler(rspTimeoutSeconds))
       pipeline
    }
 
@@ -199,6 +213,14 @@
    override def encode(ctx: ChannelHandlerContext, ch: Channel, msg: Any) = {
       trace("Encode {0} so that it's sent to the server", msg)
       msg match {
+         case partial: PartialOp => {
+            val buffer = new ChannelBufferAdapter(ChannelBuffers.dynamicBuffer)
+            buffer.writeByte(partial.magic.asInstanceOf[Byte]) // magic
+            buffer.writeUnsignedLong(partial.id) // message id
+            buffer.writeByte(10) // version
+            buffer.writeByte(partial.code) // opcode
+            buffer.getUnderlyingChannelBuffer
+         }
          case op: Op => {
             val buffer = new ChannelBufferAdapter(ChannelBuffers.dynamicBuffer)
             buffer.writeByte(op.magic.asInstanceOf[Byte]) // magic
@@ -317,7 +339,7 @@
    }
 }
 
-private class ClientHandler extends SimpleChannelUpstreamHandler {
+private class ClientHandler(rspTimeoutSeconds: Int) extends SimpleChannelUpstreamHandler {
 
    private val responses = new ConcurrentHashMap[Long, Response]
 
@@ -338,7 +360,7 @@
             i += 1
          }
       }
-      while (v == null && i < 10000)
+      while (v == null && i < (rspTimeoutSeconds * 10))
       v
    }
 
@@ -358,6 +380,20 @@
    lazy val id = HotRodClient.idCounter.incrementAndGet
 }
 
+class PartialOp(override val magic: Int,
+                override val code: Byte,
+                override val cacheName: String,
+                override val key: Array[Byte],
+                override val lifespan: Int,
+                override val maxIdle: Int,
+                override val value: Array[Byte],
+                override val flags: Int,
+                override val version: Long,
+                override val clientIntelligence: Byte,
+                override val topologyId: Int)
+      extends Op(magic, code, cacheName, key, lifespan, maxIdle, value, flags, version, clientIntelligence, topologyId) {
+}
+
 class StatsOp(override val magic: Int,
               override val code: Byte,
               override val cacheName: String,

Modified: trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/test/HotRodTestingUtil.scala
===================================================================
--- trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/test/HotRodTestingUtil.scala	2010-04-21 13:47:30 UTC (rev 1707)
+++ trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/test/HotRodTestingUtil.scala	2010-04-21 14:58:08 UTC (rev 1708)
@@ -24,13 +24,16 @@
    def startHotRodServer(manager: CacheManager): HotRodServer =
       startHotRodServer(manager, UniquePortThreadLocal.get.intValue)
 
-   def startHotRodServer(manager: CacheManager, port: Int): HotRodServer = {
+   def startHotRodServer(manager: CacheManager, port: Int): HotRodServer =
+      startHotRodServer(manager, port, 0)
+
+   def startHotRodServer(manager: CacheManager, port: Int, idleTimeout: Int): HotRodServer = {
       val server = new HotRodServer {
          override protected def defineTopologyCacheConfig(cacheManager: CacheManager) {
             // No-op since topology cache configuration comes defined by the test
          }
       }
-      server.start(host, port, manager, 0, 0)
+      server.start(host, port, manager, 0, 0, idleTimeout)
       server
    }
 

Modified: trunk/server/memcached/src/test/scala/org/infinispan/server/memcached/test/MemcachedTestingUtil.scala
===================================================================
--- trunk/server/memcached/src/test/scala/org/infinispan/server/memcached/test/MemcachedTestingUtil.scala	2010-04-21 13:47:30 UTC (rev 1707)
+++ trunk/server/memcached/src/test/scala/org/infinispan/server/memcached/test/MemcachedTestingUtil.scala	2010-04-21 14:58:08 UTC (rev 1708)
@@ -38,7 +38,7 @@
 
    def startMemcachedTextServer(cacheManager: CacheManager, port: Int): MemcachedServer = {
       val server = new MemcachedServer
-      server.start(host, port, cacheManager, 0, 0)
+      server.start(host, port, cacheManager, 0, 0, 0)
       server
    }
 
@@ -54,7 +54,7 @@
             }
          }
       }
-      server.start(host, port, cacheManager, 0, 0)
+      server.start(host, port, cacheManager, 0, 0, 0)
       server
    }
    



More information about the infinispan-commits mailing list