[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