Repository: stratos
Updated Branches:
  refs/heads/4.1.0-test 91b1af197 -> 285671553


Introducing graceful termination functionality for mock member


Project: http://git-wip-us.apache.org/repos/asf/stratos/repo
Commit: http://git-wip-us.apache.org/repos/asf/stratos/commit/28567155
Tree: http://git-wip-us.apache.org/repos/asf/stratos/tree/28567155
Diff: http://git-wip-us.apache.org/repos/asf/stratos/diff/28567155

Branch: refs/heads/4.1.0-test
Commit: 285671553ac2a7cf3078de82c2a5fd1e9c7034f7
Parents: 91b1af1
Author: Imesh Gunaratne <[email protected]>
Authored: Tue Dec 9 12:59:49 2014 +0530
Committer: Imesh Gunaratne <[email protected]>
Committed: Tue Dec 9 12:59:49 2014 +0530

----------------------------------------------------------------------
 .../controller/iaases/mock/MockConstants.java   | 27 ++++++
 .../controller/iaases/mock/MockIaasService.java |  5 +-
 .../controller/iaases/mock/MockMember.java      | 87 +++++++++++++++++---
 .../iaases/mock/MockMemberEventPublisher.java   | 21 ++---
 4 files changed, 115 insertions(+), 25 deletions(-)
----------------------------------------------------------------------


http://git-wip-us.apache.org/repos/asf/stratos/blob/28567155/components/org.apache.stratos.cloud.controller/src/main/java/org/apache/stratos/cloud/controller/iaases/mock/MockConstants.java
----------------------------------------------------------------------
diff --git 
a/components/org.apache.stratos.cloud.controller/src/main/java/org/apache/stratos/cloud/controller/iaases/mock/MockConstants.java
 
b/components/org.apache.stratos.cloud.controller/src/main/java/org/apache/stratos/cloud/controller/iaases/mock/MockConstants.java
new file mode 100644
index 0000000..cdd3490
--- /dev/null
+++ 
b/components/org.apache.stratos.cloud.controller/src/main/java/org/apache/stratos/cloud/controller/iaases/mock/MockConstants.java
@@ -0,0 +1,27 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *  http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.stratos.cloud.controller.iaases.mock;
+
+/**
+ * Mock constant definitions.
+ */
+public class MockConstants {
+    public static final int MAX_MOCK_MEMBER_COUNT = 100;
+}

http://git-wip-us.apache.org/repos/asf/stratos/blob/28567155/components/org.apache.stratos.cloud.controller/src/main/java/org/apache/stratos/cloud/controller/iaases/mock/MockIaasService.java
----------------------------------------------------------------------
diff --git 
a/components/org.apache.stratos.cloud.controller/src/main/java/org/apache/stratos/cloud/controller/iaases/mock/MockIaasService.java
 
b/components/org.apache.stratos.cloud.controller/src/main/java/org/apache/stratos/cloud/controller/iaases/mock/MockIaasService.java
index 4bf9dd4..31aa0f0 100644
--- 
a/components/org.apache.stratos.cloud.controller/src/main/java/org/apache/stratos/cloud/controller/iaases/mock/MockIaasService.java
+++ 
b/components/org.apache.stratos.cloud.controller/src/main/java/org/apache/stratos/cloud/controller/iaases/mock/MockIaasService.java
@@ -49,16 +49,17 @@ import java.util.concurrent.ExecutorService;
 public class MockIaasService {
 
     private static final Log log = LogFactory.getLog(MockIaasService.class);
+
+    private static final ExecutorService executorService = 
StratosThreadPool.getExecutorService("MOCK_MEMBER_THREAD_EXECUTOR",
+            MockConstants.MAX_MOCK_MEMBER_COUNT);
     private static final String MOCK_IAAS_MEMBERS = "/mock/iaas/members";
     private static volatile MockIaasService instance;
 
-    private ExecutorService executorService;
     private MockPartitionValidator partitionValidator;
     private ConcurrentHashMap<String, MockMember> membersMap;
 
     private MockIaasService() {
         super();
-        executorService = 
StratosThreadPool.getExecutorService("MOCK_IAAS_THREAD_EXECUTOR", 100);
         partitionValidator = new MockPartitionValidator();
         membersMap = readFromRegistry();
         if(membersMap == null) {

http://git-wip-us.apache.org/repos/asf/stratos/blob/28567155/components/org.apache.stratos.cloud.controller/src/main/java/org/apache/stratos/cloud/controller/iaases/mock/MockMember.java
----------------------------------------------------------------------
diff --git 
a/components/org.apache.stratos.cloud.controller/src/main/java/org/apache/stratos/cloud/controller/iaases/mock/MockMember.java
 
b/components/org.apache.stratos.cloud.controller/src/main/java/org/apache/stratos/cloud/controller/iaases/mock/MockMember.java
index ec4a1e4..6e4a83b 100644
--- 
a/components/org.apache.stratos.cloud.controller/src/main/java/org/apache/stratos/cloud/controller/iaases/mock/MockMember.java
+++ 
b/components/org.apache.stratos.cloud.controller/src/main/java/org/apache/stratos/cloud/controller/iaases/mock/MockMember.java
@@ -22,8 +22,15 @@ package org.apache.stratos.cloud.controller.iaases.mock;
 import org.apache.commons.logging.Log;
 import org.apache.commons.logging.LogFactory;
 import 
org.apache.stratos.cloud.controller.iaases.mock.statistics.MockHealthStatisticsNotifier;
+import org.apache.stratos.messaging.event.Event;
+import 
org.apache.stratos.messaging.event.instance.notifier.InstanceCleanupClusterEvent;
+import 
org.apache.stratos.messaging.event.instance.notifier.InstanceCleanupMemberEvent;
+import 
org.apache.stratos.messaging.listener.instance.notifier.InstanceCleanupClusterEventListener;
+import 
org.apache.stratos.messaging.listener.instance.notifier.InstanceCleanupMemberEventListener;
+import 
org.apache.stratos.messaging.message.receiver.instance.notifier.InstanceNotifierEventReceiver;
 
 import java.io.Serializable;
+import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
 import java.util.concurrent.ScheduledExecutorService;
 import java.util.concurrent.TimeUnit;
@@ -34,10 +41,13 @@ import java.util.concurrent.TimeUnit;
 public class MockMember implements Runnable, Serializable {
 
     private static final Log log = LogFactory.getLog(MockMember.class);
-    private static final ScheduledExecutorService scheduler = 
Executors.newScheduledThreadPool(1);
-    private static final int HEALTH_STAT_INTERVAL = 15;
+    private static final ExecutorService executorService =
+            Executors.newFixedThreadPool(MockConstants.MAX_MOCK_MEMBER_COUNT);
+    private static final ScheduledExecutorService scheduler =
+            
Executors.newScheduledThreadPool(MockConstants.MAX_MOCK_MEMBER_COUNT);
+    private static final int HEALTH_STAT_INTERVAL = 15; // 15 seconds
 
-    private MockMemberContext mockMemberContext;
+    private final MockMemberContext mockMemberContext;
     private boolean terminated;
 
     public MockMember(MockMemberContext mockMemberContext) {
@@ -46,7 +56,7 @@ public class MockMember implements Runnable, Serializable {
 
     @Override
     public void run() {
-        if(log.isInfoEnabled()) {
+        if (log.isInfoEnabled()) {
             log.info(String.format("Mock member started: [member-id] %s", 
mockMemberContext.getMemberId()));
         }
 
@@ -56,18 +66,69 @@ public class MockMember implements Runnable, Serializable {
         sleep(5000);
         
MockMemberEventPublisher.publishInstanceActivatedEvent(mockMemberContext);
 
-        if (log.isInfoEnabled()) {
-            log.info(String.format("Starting health statistics notifier: 
[member-id] %s", mockMemberContext.getMemberId()));
+        startInstanceNotifierReceiver();
+        startHealthStatisticsPublisher();
+
+        while (!terminated) {
+            sleep(1000);
         }
-        scheduler.scheduleAtFixedRate(new 
MockHealthStatisticsNotifier(mockMemberContext),
-                HEALTH_STAT_INTERVAL, HEALTH_STAT_INTERVAL, TimeUnit.SECONDS);
+    }
 
-        if (log.isInfoEnabled()) {
-            log.info(String.format("Health statistics notifier started: 
[member-id] %s", mockMemberContext.getMemberId()));
+    private void startInstanceNotifierReceiver() {
+        if (log.isDebugEnabled()) {
+            log.debug("Starting instance notifier event message receiver");
         }
 
-        while(!terminated) {
-            sleep(1000);
+        final InstanceNotifierEventReceiver instanceNotifierEventReceiver = 
new InstanceNotifierEventReceiver();
+        instanceNotifierEventReceiver.addEventListener(new 
InstanceCleanupClusterEventListener() {
+            @Override
+            protected void onEvent(Event event) {
+                InstanceCleanupClusterEvent instanceCleanupClusterEvent = 
(InstanceCleanupClusterEvent) event;
+                if 
(mockMemberContext.getClusterId().equals(instanceCleanupClusterEvent.getClusterId())
 &&
+                        
mockMemberContext.getInstanceId().equals(instanceCleanupClusterEvent.getInstanceId()))
 {
+                    handleMemberTermination();
+                }
+            }
+        });
+
+        instanceNotifierEventReceiver.addEventListener(new 
InstanceCleanupMemberEventListener() {
+            @Override
+            protected void onEvent(Event event) {
+                InstanceCleanupMemberEvent instanceCleanupClusterEvent = 
(InstanceCleanupMemberEvent) event;
+                if 
(mockMemberContext.getMemberId().equals(instanceCleanupClusterEvent.getMemberId()))
 {
+                    handleMemberTermination();
+                }
+            }
+        });
+
+        executorService.submit(new Runnable() {
+            @Override
+            public void run() {
+                instanceNotifierEventReceiver.execute();
+            }
+        });
+
+        if (log.isDebugEnabled()) {
+            log.debug("Instance notifier event message receiver started");
+        }
+    }
+
+    private void handleMemberTermination() {
+        
MockMemberEventPublisher.publishMaintenanceModeEvent(mockMemberContext);
+
+        sleep(2000);
+        
MockMemberEventPublisher.publishInstanceReadyToShutdownEvent(mockMemberContext);
+    }
+
+    private void startHealthStatisticsPublisher() {
+        if (log.isDebugEnabled()) {
+            log.debug(String.format("Starting health statistics notifier: 
[member-id] %s", mockMemberContext.getMemberId()));
+        }
+        scheduler.scheduleAtFixedRate(new 
MockHealthStatisticsNotifier(mockMemberContext),
+                HEALTH_STAT_INTERVAL, HEALTH_STAT_INTERVAL, TimeUnit.SECONDS);
+
+        if (log.isDebugEnabled()) {
+            log.debug(String.format("Health statistics notifier started: 
[member-id] %s", mockMemberContext.getMemberId()));
         }
     }
 
@@ -88,7 +149,7 @@ public class MockMember implements Runnable, Serializable {
         terminated = true;
         scheduler.shutdownNow();
 
-        if(log.isInfoEnabled()) {
+        if (log.isInfoEnabled()) {
             log.info(String.format("Mock member terminated: [member-id] %s", 
memberId));
         }
     }

http://git-wip-us.apache.org/repos/asf/stratos/blob/28567155/components/org.apache.stratos.cloud.controller/src/main/java/org/apache/stratos/cloud/controller/iaases/mock/MockMemberEventPublisher.java
----------------------------------------------------------------------
diff --git 
a/components/org.apache.stratos.cloud.controller/src/main/java/org/apache/stratos/cloud/controller/iaases/mock/MockMemberEventPublisher.java
 
b/components/org.apache.stratos.cloud.controller/src/main/java/org/apache/stratos/cloud/controller/iaases/mock/MockMemberEventPublisher.java
index 6d06c85..1b1ea43 100644
--- 
a/components/org.apache.stratos.cloud.controller/src/main/java/org/apache/stratos/cloud/controller/iaases/mock/MockMemberEventPublisher.java
+++ 
b/components/org.apache.stratos.cloud.controller/src/main/java/org/apache/stratos/cloud/controller/iaases/mock/MockMemberEventPublisher.java
@@ -70,8 +70,7 @@ public class MockMemberEventPublisher {
 
         // Event publisher connection will
         String topic = Util.getMessageTopicName(event);
-        EventPublisher eventPublisher = EventPublisherPool
-                .getPublisher(topic);
+        EventPublisher eventPublisher = EventPublisherPool.getPublisher(topic);
         eventPublisher.publish(event);
         if (log.isInfoEnabled()) {
             log.info("Instance activated event published");
@@ -80,7 +79,8 @@ public class MockMemberEventPublisher {
 
     public static void publishInstanceReadyToShutdownEvent(MockMemberContext 
mockMemberContext) {
         if (log.isInfoEnabled()) {
-            log.info("Publishing instance activated event");
+            log.info(String.format("Publishing instance ready to shutdown 
event: [member-id] %s",
+                    mockMemberContext.getMemberId()));
         }
         InstanceReadyToShutdownEvent event = new InstanceReadyToShutdownEvent(
                 mockMemberContext.getServiceName(),
@@ -94,14 +94,15 @@ public class MockMemberEventPublisher {
                 .getPublisher(topic);
         eventPublisher.publish(event);
         if (log.isInfoEnabled()) {
-            log.info("Instance ReadyToShutDown event published");
+            log.info(String.format("Instance ready to shutDown event 
published: [member-id] %s",
+                    mockMemberContext.getMemberId()));
         }
-
     }
 
     public static void publishMaintenanceModeEvent(MockMemberContext 
mockMemberContext) {
         if (log.isInfoEnabled()) {
-            log.info("Publishing instance maintenance mode event");
+            log.info(String.format("Publishing instance maintenance mode 
event: [member-id] %s",
+                    mockMemberContext.getMemberId()));
         }
         InstanceMaintenanceModeEvent event = new InstanceMaintenanceModeEvent(
                 mockMemberContext.getServiceName(),
@@ -111,12 +112,12 @@ public class MockMemberEventPublisher {
                 mockMemberContext.getMemberId(),
                 mockMemberContext.getInstanceId());
         String topic = Util.getMessageTopicName(event);
-        EventPublisher eventPublisher = EventPublisherPool
-                .getPublisher(topic);
+        EventPublisher eventPublisher = EventPublisherPool.getPublisher(topic);
         eventPublisher.publish(event);
+
         if (log.isInfoEnabled()) {
-            log.info("Instance Maintenance mode event published");
+            log.info(String.format("Instance Maintenance mode event published: 
[member-id] %s",
+                    mockMemberContext.getMemberId()));
         }
     }
-
 }

Reply via email to