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())); } } - }
