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

infinispan-commits at lists.jboss.org infinispan-commits at lists.jboss.org
Fri Apr 23 09:53:01 EDT 2010


Author: galder.zamarreno at jboss.com
Date: 2010-04-23 09:53:00 -0400 (Fri, 23 Apr 2010)
New Revision: 1724

Modified:
   trunk/server/
   trunk/server/core/src/main/scala/org/infinispan/server/core/AbstractProtocolServer.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/main/scala/org/infinispan/server/hotrod/TopologyAddress.scala
   trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/HotRodReplicationTest.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
Log:
[ISPN-384] (Implement topology and hash distribution headers in Hot Rod) Added logic to update Hot Rod topology info when members have crashed or stopped responding and they've been evicted from the JGroups view.


Property changes on: trunk/server
___________________________________________________________________
Name: svn:ignore
   + target
*.log


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-23 13:07:09 UTC (rev 1723)
+++ trunk/server/core/src/main/scala/org/infinispan/server/core/AbstractProtocolServer.scala	2010-04-23 13:53:00 UTC (rev 1724)
@@ -2,7 +2,7 @@
 
 import java.net.InetSocketAddress
 import transport.netty.{EncoderAdapter, NettyTransport}
-import transport.{Decoder, Encoder, Transport}
+import transport.Transport
 import org.infinispan.manager.CacheManager
 import org.infinispan.server.core.VersionGenerator._
 
@@ -41,6 +41,8 @@
 
    def getCacheManager = cacheManager
 
+   def getHost = host
+
    def getPort = port
 
 }

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-23 13:07:09 UTC (rev 1723)
+++ trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/NettyTransport.scala	2010-04-23 13:53:00 UTC (rev 1724)
@@ -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.{HashedWheelTimer, ThreadNameDeterminer, ThreadRenamingRunnable}
+import org.jboss.netty.util.{ThreadNameDeterminer, ThreadRenamingRunnable}
 
 /**
  * // TODO: Document this
@@ -97,14 +97,4 @@
 
 }
 
-object NettyTransport extends Logging
-
-private class NamedThreadFactory(val name: String) extends ThreadFactory {
-   val threadCounter = new AtomicInteger
-
-   override def newThread(r: Runnable): Thread = {
-      var t = new Thread(r, System.getProperty("program.name") + "-" + name + '-' + threadCounter.incrementAndGet)
-      t.setDaemon(true)
-      t
-   }
-}
+object NettyTransport extends Logging
\ No newline at end of file

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-23 13:07:09 UTC (rev 1723)
+++ trunk/server/hotrod/src/main/scala/org/infinispan/server/hotrod/HotRodServer.scala	2010-04-23 13:53:00 UTC (rev 1724)
@@ -2,11 +2,17 @@
 
 import org.infinispan.manager.CacheManager
 import org.infinispan.server.core.transport.{Decoder, Encoder}
-import org.jgroups.blocks.RequestOptions
 import org.infinispan.server.core.{Logging, AbstractProtocolServer}
 import org.infinispan.config.Configuration
 import org.infinispan.config.Configuration.CacheMode
 import org.infinispan.Cache
+import org.infinispan.notifications.Listener
+import org.infinispan.notifications.cachemanagerlistener.annotation.ViewChanged
+import org.infinispan.notifications.cachemanagerlistener.event.ViewChangedEvent
+import scala.collection.JavaConversions._
+import org.infinispan.remoting.transport.Address
+import java.util.concurrent.{ThreadFactory, Callable, Executors}
+import java.util.concurrent.atomic.AtomicInteger
 
 /**
  * // TODO: Document this
@@ -15,40 +21,73 @@
  */
 
 class HotRodServer extends AbstractProtocolServer("HotRod") {
-
    import HotRodServer._
+   private var isClustered: Boolean = _
+   private var address: TopologyAddress = _
 
+   def getAddress: TopologyAddress = address
+
    override def getEncoder: Encoder = new HotRodEncoder
 
    override def getDecoder: Decoder = new HotRodDecoder(getCacheManager)
 
    override def start(host: String, port: Int, cacheManager: CacheManager, masterThreads: Int, workerThreads: Int, idleTimeout: Int) {
       super.start(host, port, cacheManager, masterThreads, workerThreads, idleTimeout)
+      isClustered = cacheManager.getGlobalConfiguration.getTransportClass != null
       // If clustered, set up a cache for topology information
-      if (cacheManager.getGlobalConfiguration.getTransportClass != null) {
-         defineTopologyCacheConfig(cacheManager)
-         val topologyCache: Cache[String, TopologyView] = cacheManager.getCache(TopologyCacheName)
-         val currentView = topologyCache.get("view")
-         if (currentView != null) {
-            // TODO: If distribution configured, add hashcode of this address
-            val newMembers = currentView.members ::: List(TopologyAddress(host, port, 0))
-            val newView = TopologyView(currentView.topologyId + 1, newMembers)
-            val replaced = topologyCache.replace("view", currentView, newView)
-            if (!replaced) {
-               // TODO: There was a concurrent view modification, get and try to install new view again.
-            }
+      if (isClustered)
+         addSelfToTopologyView(host, port, cacheManager)
+   }
+
+   private def addSelfToTopologyView(host: String, port: Int, cacheManager: CacheManager) {
+      defineTopologyCacheConfig(cacheManager)
+      val topologyCache: Cache[String, TopologyView] = cacheManager.getCache(TopologyCacheName)
+      cacheManager.addListener(new CrashedMemberDetectorListener)
+      address = TopologyAddress(host, port, 0, cacheManager.getAddress)
+      val currentView = topologyCache.get("view")
+      // TODO: If distribution configured, add hashcode of this address
+      if (currentView != null) {
+         val newMembers = currentView.members ::: List(address)
+         val newView = TopologyView(currentView.topologyId + 1, newMembers)
+         val replaced = topologyCache.replace("view", currentView, newView)
+         if (!replaced) {
+            // TODO: There was a concurrent view modification, get and try to install new view again.
          } else {
-            // TODO add check for distribution and if so, put the right hashcode
-            val newMembers = List(TopologyAddress(host, port, 0))
-            val newView = TopologyView(1, newMembers)
-            val prev = topologyCache.putIfAbsent("view", newView)
-            if (prev != null) {
-               // TODO: There was a concurrent view modification, get and try to install new view again.
-            }
+            debug("Added {0} to topology, new view is {1}", address, newView)
          }
+      } else {
+         val newMembers = List(address)
+         val newView = TopologyView(1, newMembers)
+         val prev = topologyCache.putIfAbsent("view", newView)
+         if (prev != null) {
+            // TODO: There was a concurrent view modification, get and try to install new view again.
+         } else {
+            debug("First member to start, topology view is {0}", newView)
+         }
       }
    }
 
+   override def stop {
+      super.stop
+      if (isClustered)
+         removeSelfFromTopologyView
+   }
+
+   protected def removeSelfFromTopologyView {
+      // Graceful shutdown, remove this node as member and install new view
+      val topologyCache: Cache[String, TopologyView] = getCacheManager.getCache(TopologyCacheName)
+      val currentView = topologyCache.get("view")
+      // TODO: If distribution configured, add hashcode of this address
+      val newMembers = currentView.members.filterNot(_ == address)
+      val newView = TopologyView(currentView.topologyId + 1, newMembers)
+      val replaced = topologyCache.replace("view", currentView, newView)
+      if (!replaced) {
+         // TODO: There was a concurrent view modification. Just give up, logic to deal with crashed/stalled members will deal with this
+      } else {
+         debug("Removed {0} from topology view, new view is {1}", address, newView)
+      }
+   }
+
    protected def defineTopologyCacheConfig(cacheManager: CacheManager) {
       val topologyCacheConfig = new Configuration
       topologyCacheConfig.setCacheMode(CacheMode.REPL_SYNC)
@@ -57,8 +96,82 @@
       cacheManager.defineConfiguration(TopologyCacheName, topologyCacheConfig)
    }
 
+   @Listener
+   class CrashedMemberDetectorListener {
+      import HotRodServer._
+
+      private val executor = Executors.newCachedThreadPool(new ThreadFactory(){
+         val threadCounter = new AtomicInteger
+
+         override def newThread(r: Runnable): Thread = {
+            var t = new Thread(r, "CrashedMemberDetectorThread-" + threadCounter.incrementAndGet)
+            t.setDaemon(true)
+            t
+         }
+      })
+
+      @ViewChanged
+      def handleViewChange(e: ViewChangedEvent) {
+         val cacheManager = e.getCacheManager
+         // Only the coordinator can potentially make modifications related to crashed members.
+         // This is to avoid all nodes trying to make the same modification which would be wasteful and lead to deadlocks.
+         if (cacheManager.isCoordinator) {
+            // Use a separate thread to avoid blocking the view handler thread
+            val callable = new Callable[Void] {
+               override def call = {
+                  try {
+                     val newMembers = e.getNewMembers
+                     val oldMembers = e.getOldMembers
+                     // Someone left the cluster, verify whether it did it gracefully or crashed.
+                     if (oldMembers.size > newMembers.size) {
+                        val topologyCache: Cache[String, TopologyView] = getCacheManager.getCache(TopologyCacheName)
+                        val currentView = topologyCache.get("view")
+                        var newTopologyMembers = currentView.members
+                        for (oldMember <- asIterator(oldMembers.iterator)) {
+                           // If old member is not amongst the new ones, check whether it's still in the topology cache
+                           if (!newMembers.contains(oldMember)) {
+                              trace("Old member {0} is not in new view {1}, did it crash?", oldMember, newMembers)
+                              // If old memmber is in topology, it means that it had an abnormal ending
+                              val (isCrashed, crashedTopologyMember) = isOldMemberInTopology(oldMember, currentView)
+                              if (isCrashed) {
+                                 trace("Old member {0} with topology address {1} is still present in Hot Rod topology " +
+                                       "{2}, so must have crashed.", oldMember, crashedTopologyMember, currentView)
+                                 newTopologyMembers = newTopologyMembers.filterNot(_ == crashedTopologyMember)
+                                 trace("After removal, new Hot Rod topology is {0}", newTopologyMembers)
+                              }
+                           }
+                        }
+                        if (newTopologyMembers.size < currentView.members.size) {
+                           val newView = TopologyView(currentView.topologyId + 1, newTopologyMembers)
+                           val replaced = topologyCache.replace("view", currentView, newView)
+                           if (!replaced) {
+                              // TODO: How to deal with concurrent failures at this point?
+                           }
+                        }
+                     }
+                  } catch {
+                     case t: Throwable => error("Error detecting crashed member", t)
+                  }
+                  null
+               }
+            }
+            executor.submit(callable);
+         }
+      }
+
+      private def isOldMemberInTopology(oldMember: Address, view: TopologyView): (Boolean, TopologyAddress) = {
+         for (member <- view.members) {
+            if (member.clusterAddress == oldMember) {
+               return (true, member)
+            }
+         }
+         (false, null)
+      }
+   }
+
 }
 
-object HotRodServer {
+object HotRodServer extends Logging {
    val TopologyCacheName = "___hotRodTopologyCache"
-}
\ No newline at end of file
+}
+

Modified: trunk/server/hotrod/src/main/scala/org/infinispan/server/hotrod/TopologyAddress.scala
===================================================================
--- trunk/server/hotrod/src/main/scala/org/infinispan/server/hotrod/TopologyAddress.scala	2010-04-23 13:07:09 UTC (rev 1723)
+++ trunk/server/hotrod/src/main/scala/org/infinispan/server/hotrod/TopologyAddress.scala	2010-04-23 13:53:00 UTC (rev 1724)
@@ -2,6 +2,7 @@
 
 import java.io.{ObjectInput, ObjectOutput}
 import org.infinispan.marshall.Marshallable
+import org.infinispan.remoting.transport.Address
 
 /**
  * // TODO: Document this
@@ -9,7 +10,7 @@
  * @since // TODO
  */
 @Marshallable(externalizer = classOf[TopologyAddress.Externalizer], id = 58)
-case class TopologyAddress(val host: String, val port: Int, val hostHashCode: Int)
+case class TopologyAddress(val host: String, val port: Int, val hostHashCode: Int, val clusterAddress: Address)
 
 object TopologyAddress {
    class Externalizer extends org.infinispan.marshall.Externalizer {
@@ -18,13 +19,15 @@
          output.writeObject(topologyAddress.host)
          output.writeInt(topologyAddress.port)
          output.writeInt(topologyAddress.hostHashCode)
+         output.writeObject(topologyAddress.clusterAddress)
       }
 
       override def readObject(input: ObjectInput): AnyRef = {
          val host = input.readObject.asInstanceOf[String]
          val port = input.readInt
          val hostHashCode = input.readInt
-         TopologyAddress(host, port, hostHashCode)
+         val clusterAddress = input.readObject.asInstanceOf[Address]
+         TopologyAddress(host, port, hostHashCode, clusterAddress)
       }
    }
 }
\ 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-23 13:07:09 UTC (rev 1723)
+++ trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/HotRodReplicationTest.scala	2010-04-23 13:53:00 UTC (rev 1724)
@@ -1,6 +1,5 @@
 package org.infinispan.server.hotrod
 
-import org.infinispan.test.MultipleCacheManagersTest
 import org.infinispan.config.Configuration
 import java.lang.reflect.Method
 import test.HotRodClient
@@ -9,6 +8,7 @@
 import org.infinispan.config.Configuration.CacheMode
 import org.testng.Assert._
 import org.testng.annotations.{AfterMethod, AfterClass, Test}
+import org.infinispan.test.{TestingUtil, MultipleCacheManagersTest}
 
 /**
  * // TODO: Document this
@@ -41,10 +41,10 @@
 
    @AfterClass(alwaysRun = true)
    override def destroy {
-      super.destroy
       log.debug("Test finished, close Hot Rod server", null)
       clients.foreach(_.stop)
       servers.foreach(_.stop)
+      super.destroy // Stop the caches last so that at stoppage time topology cache can be updated properly
    }
 
    @AfterMethod(alwaysRun=true)
@@ -118,11 +118,11 @@
    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", servers.head.getPort, 0))
-      assertEquals(topologyResp.view.members.tail.head, TopologyAddress("127.0.0.1", servers.tail.head.getPort, 0))
+      assertAddressEquals(topologyResp.view.members.head, servers.head.getAddress)
+      assertAddressEquals(topologyResp.view.members.tail.head, servers.tail.head.getAddress)
    }
 
-   def testReplicatedPutWithTopologyAwareClient(m: Method) {
+   def testReplicatedPutWithTopologyChanges(m: Method) {
       var resp = clients.head.put(k(m) , 0, 0, v(m), 1, 0)
       assertStatus(resp.status, Success)
       assertEquals(resp.topologyResponse, None)
@@ -138,22 +138,65 @@
       assertEquals(resp.topologyResponse, None)
       assertSuccess(clients.tail.head.get(k(m), 0), v(m, "v3-"))
 
-      val cm = addClusterEnabledCacheManager()
+      var cm = addClusterEnabledCacheManager()
       cm.defineConfiguration(cacheName, createCacheConfig)
       cm.defineConfiguration(TopologyCacheName, createTopologyCacheConfig)
-      servers = servers ::: List(startHotRodServer(cacheManagers.get(2), servers.tail.head.getPort + 25)) 
+      val newServer = startHotRodServer(cm, servers.tail.head.getPort + 25)
+      servers = servers ::: List(newServer)
 
       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))
+      assertAddressEquals(resp.topologyResponse.get.view.members.head, servers.head.getAddress)
+      assertAddressEquals(resp.topologyResponse.get.view.members.tail.head, servers.tail.head.getAddress)
+      assertAddressEquals(resp.topologyResponse.get.view.members.tail.tail.head, servers.tail.tail.head.getAddress)
       assertSuccess(clients.tail.head.get(k(m), 0), v(m, "v4-"))
 
-//      // TODO: Add stopping a server
-//      servers.tail.tail.head.stop
+      servers.tail.tail.head.stop
+      servers = servers.filterNot(_ == newServer)
+      cm.stop
+
+      resp = clients.head.put(k(m) , 0, 0, v(m, "v5-"), 2, 3)
+      assertStatus(resp.status, Success)
+      assertEquals(resp.topologyResponse.get.view.topologyId, 4)
+      assertEquals(resp.topologyResponse.get.view.members.size, 2)
+      assertAddressEquals(resp.topologyResponse.get.view.members.head, servers.head.getAddress)
+      assertAddressEquals(resp.topologyResponse.get.view.members.tail.head, servers.tail.head.getAddress)
+      assertSuccess(clients.tail.head.get(k(m), 0), v(m, "v5-"))
+
+      cm = addClusterEnabledCacheManager()
+      cm.defineConfiguration(cacheName, createCacheConfig)
+      cm.defineConfiguration(TopologyCacheName, createTopologyCacheConfig)
+      val crashingServer = startCrashingHotRodServer(cm, servers.tail.head.getPort + 11)
+      servers = servers ::: List(crashingServer)
+
+      resp = clients.head.put(k(m) , 0, 0, v(m, "v6-"), 2, 4)
+      assertStatus(resp.status, Success)
+      assertEquals(resp.topologyResponse.get.view.topologyId, 5)
+      assertEquals(resp.topologyResponse.get.view.members.size, 3)
+      assertAddressEquals(resp.topologyResponse.get.view.members.head, servers.head.getAddress)
+      assertAddressEquals(resp.topologyResponse.get.view.members.tail.head, servers.tail.head.getAddress)
+      assertAddressEquals(resp.topologyResponse.get.view.members.tail.tail.head, servers.tail.tail.head.getAddress)
+      assertSuccess(clients.tail.head.get(k(m), 0), v(m, "v6-"))
+
+      crashingServer.stop
+      servers = servers.filterNot(_ == crashingServer)
+      cm.stop
+      TestingUtil.blockUntilViewsReceived(10000, true, manager(0), manager(1))
+
+      resp = clients.head.put(k(m) , 0, 0, v(m, "v7-"), 2, 5)
+      assertStatus(resp.status, Success)
+      assertEquals(resp.topologyResponse.get.view.topologyId, 6)
+      assertEquals(resp.topologyResponse.get.view.members.size, 2)
+      assertAddressEquals(resp.topologyResponse.get.view.members.head, servers.head.getAddress)
+      assertAddressEquals(resp.topologyResponse.get.view.members.tail.head, servers.tail.head.getAddress)
+      assertSuccess(clients.tail.head.get(k(m), 0), v(m, "v7-"))
    }
 
+   private def assertAddressEquals(actual: TopologyAddress, expected: TopologyAddress) {
+      assertEquals(actual.host, expected.host)
+      assertEquals(actual.port, expected.port)
+      assertEquals(actual.hostHashCode, expected.hostHashCode)
+   }
 }
\ 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-23 13:07:09 UTC (rev 1723)
+++ trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/test/HotRodClient.scala	2010-04-23 13:53:00 UTC (rev 1724)
@@ -277,7 +277,7 @@
                for (i <- 0 until numberClusterMembers) {
                   val host = buf.readString
                   val port = buf.readUnsignedShort
-                  viewArray(i) = TopologyAddress(host, port, 0)
+                  viewArray(i) = TopologyAddress(host, port, 0, null)
                }
                Some(TopologyAwareResponse(TopologyView(topologyId, viewArray.toList)))
             } else if (op.clientIntelligence == 3) {

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-23 13:07:09 UTC (rev 1723)
+++ trunk/server/hotrod/src/test/scala/org/infinispan/server/hotrod/test/HotRodTestingUtil.scala	2010-04-23 13:53:00 UTC (rev 1724)
@@ -37,6 +37,21 @@
       server
    }
 
+   def startCrashingHotRodServer(manager: CacheManager, port: Int): HotRodServer = {
+      val server = new HotRodServer {
+         override protected def defineTopologyCacheConfig(cacheManager: CacheManager) {
+            // No-op since topology cache configuration comes defined by the test
+         }
+
+         override protected def removeSelfFromTopologyView {
+            // Empty to emulate a member that's crashed/unresponsive and has not executed removal,
+            // but has been evicted from JGroups cluster.
+         }
+      }
+      server.start(host, port, manager, 0, 0, 0)
+      server
+   }
+
    def k(m: Method, prefix: String): Array[Byte] = {
       val bytes: Array[Byte] = (prefix + m.getName).getBytes
       trace("String {0} is converted to {1} bytes", prefix + m.getName, Util.printArray(bytes, true))



More information about the infinispan-commits mailing list