Repository: stratos Updated Branches: refs/heads/master 95d1faedc -> 8654829fa
fixing JCA 4.1 Project: http://git-wip-us.apache.org/repos/asf/stratos/repo Commit: http://git-wip-us.apache.org/repos/asf/stratos/commit/8654829f Tree: http://git-wip-us.apache.org/repos/asf/stratos/tree/8654829f Diff: http://git-wip-us.apache.org/repos/asf/stratos/diff/8654829f Branch: refs/heads/master Commit: 8654829fa3afc222593e9d13c5d8e0378b95bba7 Parents: 95d1fae Author: Martin Eppel <[email protected]> Authored: Fri Jan 30 10:10:21 2015 -0800 Committer: Martin Eppel <[email protected]> Committed: Fri Jan 30 10:10:21 2015 -0800 ---------------------------------------------------------------------- .../stratos/cartridge/agent/CartridgeAgent.java | 425 ++++--------------- .../agent/CartridgeAgentEventListeners.java | 418 ++++++++++++++++++ .../config/CartridgeAgentConfiguration.java | 50 ++- .../publisher/CartridgeAgentEventPublisher.java | 27 ++ .../agent/util/CartridgeAgentConstants.java | 3 + 5 files changed, 591 insertions(+), 332 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/stratos/blob/8654829f/components/org.apache.stratos.cartridge.agent/src/main/java/org/apache/stratos/cartridge/agent/CartridgeAgent.java ---------------------------------------------------------------------- diff --git a/components/org.apache.stratos.cartridge.agent/src/main/java/org/apache/stratos/cartridge/agent/CartridgeAgent.java b/components/org.apache.stratos.cartridge.agent/src/main/java/org/apache/stratos/cartridge/agent/CartridgeAgent.java index ffaaad4..946f2f9 100644 --- a/components/org.apache.stratos.cartridge.agent/src/main/java/org/apache/stratos/cartridge/agent/CartridgeAgent.java +++ b/components/org.apache.stratos.cartridge.agent/src/main/java/org/apache/stratos/cartridge/agent/CartridgeAgent.java @@ -67,6 +67,8 @@ public class CartridgeAgent implements Runnable { private static final Log log = LogFactory.getLog(CartridgeAgent.class); private static final ExtensionHandler extensionHandler = new DefaultExtensionHandler(); private boolean terminated; + + private CartridgeAgentEventListeners eventListenerns; // We have an asynchronous activity running to respond to ADC updates. We want to ensure // that no publishInstanceActivatedEvent() call is made *before* the port activation test @@ -78,19 +80,42 @@ public class CartridgeAgent implements Runnable { if (log.isInfoEnabled()) { log.info("Cartridge agent started"); } + + eventListenerns = new CartridgeAgentEventListeners(); + + log.debug("MArtin: before validating system properties"); validateRequiredSystemProperties(); + + log.debug("MArtin: after validating system properties"); + + if (log.isInfoEnabled()) { + log.info("Cartridge agent validated system properties done"); + } // Start instance notifier listener thread portsActivated = false; subscribeToTopicsAndRegisterListeners(); + + if (log.isInfoEnabled()) { + log.info("Cartridge agent subscribeToTopicsAndRegisterListeners done"); + } // Start topology event receiver thread registerTopologyEventListeners(); + if (log.isInfoEnabled()) { + log.info("Cartridge agent registerTopologyEventListeners done"); + } + // Start tenant event receiver thread registerTenantEventListeners(); + + if (log.isInfoEnabled()) { + log.info("Cartridge agent registering all event listeners ... done"); + } + // wait till the member spawned event while (!CartridgeAgentConfiguration.getInstance().isInitialized()) { try { @@ -101,9 +126,17 @@ public class CartridgeAgent implements Runnable { } catch (InterruptedException ignore) { } } + + if (log.isInfoEnabled()) { + log.info("Cartridge agent initialized done"); + } // Execute instance started shell script extensionHandler.onInstanceStartedEvent(); + + if (log.isInfoEnabled()) { + log.info("Cartridge agent onInstanceStartedEvent done"); + } // Publish instance started event CartridgeAgentEventPublisher.publishInstanceStartedEvent(); @@ -116,14 +149,26 @@ public class CartridgeAgent implements Runnable { log.error("Error processing start servers event", e); } } + + if (log.isInfoEnabled()) { + log.info("Cartridge agent startServerExtension done"); + } // Wait for all ports to be active CartridgeAgentUtils.waitUntilPortsActive(CartridgeAgentConfiguration.getInstance().getListenAddress(), CartridgeAgentConfiguration.getInstance().getPorts()); portsActivated = true; + + if (log.isInfoEnabled()) { + log.info("Cartridge agent portsActivated done"); + } // Publish instance activated event CartridgeAgentEventPublisher.publishInstanceActivatedEvent(); + + if (log.isInfoEnabled()) { + log.info("Cartridge agent publishInstanceActivatedEvent done"); + } // Check repo url String repoUrl = CartridgeAgentConfiguration.getInstance().getRepoUrl(); @@ -146,6 +191,10 @@ public class CartridgeAgent implements Runnable { // ), // 0, 10, TimeUnit.SECONDS); } */ + + if (log.isInfoEnabled()) { + log.info("Cartridge agent getRepoUrl done"); + } if ("null".equals(repoUrl) || StringUtils.isBlank(repoUrl)) { if (log.isInfoEnabled()) { @@ -153,10 +202,17 @@ public class CartridgeAgent implements Runnable { } // Execute instance activated shell script extensionHandler.onInstanceActivatedEvent(); + + if (log.isInfoEnabled()) { + log.info("Cartridge agent onInstanceActivatedEvent done"); + } // Publish instance activated event CartridgeAgentEventPublisher.publishInstanceActivatedEvent(); } else { + if (log.isInfoEnabled()) { + log.info("Cartridge agent - artifact repository found"); + } //Start periodical file processor task /*if (CartridgeAgentConfiguration.getInstance().isCommitsEnabled()) { log.info(" Commits enabled. Starting File listener "); @@ -179,7 +235,7 @@ public class CartridgeAgent implements Runnable { // CartridgeAgentConstants.SUPERTENANT_TEMP_PATH), 0, // 10, TimeUnit.SECONDS); // } - + String persistenceMappingsPayload = CartridgeAgentConfiguration.getInstance().getPersistenceMappings(); if (persistenceMappingsPayload != null) { extensionHandler.volumeMountExtension(persistenceMappingsPayload); @@ -188,7 +244,7 @@ public class CartridgeAgent implements Runnable { // start log publishing LogPublisherManager logPublisherManager = new LogPublisherManager(); publishLogs(logPublisherManager); - + // Keep the thread live until terminated while (!terminated) { try { @@ -201,372 +257,79 @@ public class CartridgeAgent implements Runnable { } protected void subscribeToTopicsAndRegisterListeners() { - if (log.isDebugEnabled()) { - log.debug("Starting instance notifier event message receiver thread"); + if (log.isDebugEnabled()) { + log.debug("SsubscribeToTopicsAndRegisterListeners before"); } - - InstanceNotifierEventReceiver instanceNotifierEventReceiver = new InstanceNotifierEventReceiver(); - instanceNotifierEventReceiver.addEventListener(new ArtifactUpdateEventListener() { - @Override - protected void onEvent(Event event) { - try { - extensionHandler.onArtifactUpdatedEvent((ArtifactUpdatedEvent) event); - } catch (Exception e) { - if (log.isErrorEnabled()) { - log.error("Error processing artifact update event", e); - } - } - } - }); - - instanceNotifierEventReceiver.addEventListener(new InstanceCleanupMemberEventListener() { - @Override - protected void onEvent(Event event) { - try { - String memberIdInPayload = CartridgeAgentConfiguration.getInstance().getMemberId(); - InstanceCleanupMemberEvent instanceCleanupMemberEvent = (InstanceCleanupMemberEvent) event; - if (memberIdInPayload.equals(instanceCleanupMemberEvent.getMemberId())) { - extensionHandler.onInstanceCleanupMemberEvent(instanceCleanupMemberEvent); - } - } catch (Exception e) { - if (log.isErrorEnabled()) { - log.error("Error processing instance cleanup member event", e); - } - } - - } - }); - - instanceNotifierEventReceiver.addEventListener(new InstanceCleanupClusterEventListener() { - @Override - protected void onEvent(Event event) { - String clusterIdInPayload = CartridgeAgentConfiguration.getInstance().getClusterId(); - InstanceCleanupClusterEvent instanceCleanupClusterEvent = (InstanceCleanupClusterEvent) event; - if (clusterIdInPayload.equals(instanceCleanupClusterEvent.getClusterId())) { - extensionHandler.onInstanceCleanupClusterEvent(instanceCleanupClusterEvent); - } - } - }); - - instanceNotifierEventReceiver.execute(); - - if(log.isInfoEnabled()) { - log.info("Instance notifier event message receiver thread started"); - } - - // Wait until message receiver is subscribed to the topic to send the instance started event - while (!instanceNotifierEventReceiver.isSubscribed()) { - try { - Thread.sleep(2000); - } catch (InterruptedException e) { - } + + eventListenerns.startInstanceNotifierReceiver(); + + if (log.isDebugEnabled()) { + log.debug("SsubscribeToTopicsAndRegisterListeners after"); } } protected void registerTopologyEventListeners() { - if (log.isDebugEnabled()) { - log.debug("Starting topology event message receiver thread"); + if (log.isDebugEnabled()) { + log.debug("registerTopologyEventListeners before"); } - - TopologyEventReceiver topologyEventReceiver = new TopologyEventReceiver(); - - topologyEventReceiver.addEventListener(new MemberCreatedEventListener() { - @Override - protected void onEvent(Event event) { - try { - boolean initialized = CartridgeAgentConfiguration.getInstance().isInitialized(); - if (initialized) { - // no need to process this event, if the member is initialized. - return; - } - TopologyManager.acquireReadLock(); - if (log.isDebugEnabled()) { - log.debug("Member created event received"); - } - MemberCreatedEvent memberCreatedEvent = (MemberCreatedEvent) event; - extensionHandler.onMemberCreatedEvent(memberCreatedEvent); - } catch (Exception e) { - if (log.isErrorEnabled()) { - log.error("Error processing member created event", e); - } - } finally { - TopologyManager.releaseReadLock(); - } - } - }); - - topologyEventReceiver.addEventListener(new MemberActivatedEventListener() { - @Override - protected void onEvent(Event event) { - boolean initialized = CartridgeAgentConfiguration.getInstance().isInitialized(); - if (!initialized) { - return; - } - try { - TopologyManager.acquireReadLock(); - if (log.isDebugEnabled()) { - log.debug("Member activated event received"); - } - MemberActivatedEvent memberActivatedEvent = (MemberActivatedEvent) event; - extensionHandler.onMemberActivatedEvent(memberActivatedEvent); - } catch (Exception e) { - if (log.isErrorEnabled()) { - log.error("Error processing member activated event", e); - } - } finally { - TopologyManager.releaseReadLock(); - } - } - }); - - topologyEventReceiver.addEventListener(new MemberTerminatedEventListener() { - @Override - protected void onEvent(Event event) { - boolean initialized = CartridgeAgentConfiguration.getInstance().isInitialized(); - if (!initialized) { - return; - } - try { - TopologyManager.acquireReadLock(); - if (log.isDebugEnabled()) { - log.debug("Member terminated event received"); - } - MemberTerminatedEvent memberTerminatedEvent = (MemberTerminatedEvent) event; - extensionHandler.onMemberTerminatedEvent(memberTerminatedEvent); - } catch (Exception e) { - if (log.isErrorEnabled()) { - log.error("Error processing member terminated event", e); - } - } finally { - TopologyManager.releaseReadLock(); - } - } - }); - - topologyEventReceiver.addEventListener(new MemberSuspendedEventListener() { - @Override - protected void onEvent(Event event) { - boolean initialized = CartridgeAgentConfiguration.getInstance().isInitialized(); - if (!initialized) { - return; - } - try { - TopologyManager.acquireReadLock(); - if (log.isDebugEnabled()) { - log.debug("Member suspended event received"); - } - MemberSuspendedEvent memberSuspendedEvent = (MemberSuspendedEvent) event; - extensionHandler.onMemberSuspendedEvent(memberSuspendedEvent); - } catch (Exception e) { - if (log.isErrorEnabled()) { - log.error("Error processing member suspended event", e); - } - } finally { - TopologyManager.releaseReadLock(); - } - } - }); - - topologyEventReceiver.addEventListener(new CompleteTopologyEventListener() { - - @Override - protected void onEvent(Event event) { - boolean initialized = CartridgeAgentConfiguration.getInstance().isInitialized(); - if (!initialized) { - try { - TopologyManager.acquireReadLock(); - if (log.isDebugEnabled()) { - log.debug("Complete topology event received"); - } - CompleteTopologyEvent completeTopologyEvent = (CompleteTopologyEvent) event; - extensionHandler.onCompleteTopologyEvent(completeTopologyEvent); - } catch (Exception e) { - if (log.isErrorEnabled()) { - log.error("Error processing complete topology event", e); - } - } finally { - TopologyManager.releaseReadLock(); - } - } - } - }); - - topologyEventReceiver.addEventListener(new MemberStartedEventListener() { - @Override - protected void onEvent(Event event) { - boolean initialized = CartridgeAgentConfiguration.getInstance().isInitialized(); - if (!initialized) { - return; - } - try { - TopologyManager.acquireReadLock(); - if (log.isDebugEnabled()) { - log.debug("Member started event received"); - } - MemberStartedEvent memberStartedEvent = (MemberStartedEvent) event; - extensionHandler.onMemberStartedEvent(memberStartedEvent); - } catch (Exception e) { - if (log.isErrorEnabled()) { - log.error("Error processing member started event", e); - } - } finally { - TopologyManager.releaseReadLock(); - } - } - }); - - topologyEventReceiver.execute(); - - if (log.isDebugEnabled()) { - log.info("Cartridge Agent topology receiver thread started"); + eventListenerns.startTopologyEventReceiver(); + + if (log.isDebugEnabled()) { + log.debug("registerTopologyEventListeners after"); } } protected void registerTenantEventListeners() { - - if (log.isDebugEnabled()) { - log.debug("Starting tenant event message receiver thread"); + if (log.isDebugEnabled()) { + log.debug("registerTenantEventListeners before X"); } - TenantEventReceiver tenantEventReceiver = new TenantEventReceiver(); - tenantEventReceiver.addEventListener(new DomainMappingAddedEventListener() { - @Override - protected void onEvent(Event event) { - try { - TenantManager.acquireReadLock(); - if (log.isDebugEnabled()) { - log.debug("Subscription domain added event received"); - } - DomainMappingAddedEvent subscriptionDomainAddedEvent = (DomainMappingAddedEvent) event; - extensionHandler.onSubscriptionDomainAddedEvent(subscriptionDomainAddedEvent); - } catch (Exception e) { - if (log.isErrorEnabled()) { - log.error("Error processing subscription domains added event", e); - } - } finally { - TenantManager.releaseReadLock(); - } - - } - }); - - tenantEventReceiver.addEventListener(new DomainMappingRemovedEventListener() { - @Override - protected void onEvent(Event event) { - try { - TenantManager.acquireReadLock(); - if (log.isDebugEnabled()) { - log.debug("Subscription domain removed event received"); - } - DomainMappingRemovedEvent subscriptionDomainRemovedEvent = (DomainMappingRemovedEvent) event; - extensionHandler.onSubscriptionDomainRemovedEvent(subscriptionDomainRemovedEvent); - } catch (Exception e) { - if (log.isErrorEnabled()) { - log.error("Error processing subscription domains removed event", e); - } - } finally { - TenantManager.releaseReadLock(); - } - } - }); - - tenantEventReceiver.addEventListener(new CompleteTenantEventListener() { - private boolean initialized; - @Override - protected void onEvent(Event event) { - if (!initialized) { - try { - TenantManager.acquireReadLock(); - if (log.isDebugEnabled()) { - log.debug("Complete tenant event received"); - } - CompleteTenantEvent completeTenantEvent = (CompleteTenantEvent) event; - extensionHandler.onCompleteTenantEvent(completeTenantEvent); - initialized = true; - } catch (Exception e) { - if (log.isErrorEnabled()) { - log.error("Error processing complete tenant event", e); - } - } finally { - TenantManager.releaseReadLock(); - } - - } else { - if (log.isInfoEnabled()) { - log.info("Complete tenant event updating task disabled"); - } - } - } - }); - - tenantEventReceiver.addEventListener(new TenantSubscribedEventListener() { - @Override - protected void onEvent(Event event) { - try { - TenantManager.acquireReadLock(); - if (log.isDebugEnabled()) { - log.debug("Tenant subscribed event received"); - } - TenantSubscribedEvent tenantSubscribedEvent = (TenantSubscribedEvent) event; - extensionHandler.onTenantSubscribedEvent(tenantSubscribedEvent); - } catch (Exception e) { - if (log.isErrorEnabled()) { - log.error("Error processing tenant subscribed event", e); - } - } finally { - TenantManager.releaseReadLock(); - } - } - }); - - tenantEventReceiver.addEventListener(new TenantUnSubscribedEventListener() { - @Override - protected void onEvent(Event event) { - try { - TenantManager.acquireReadLock(); - if (log.isDebugEnabled()) { - log.debug("Tenant unSubscribed event received"); - } - TenantUnSubscribedEvent tenantUnSubscribedEvent = (TenantUnSubscribedEvent) event; - extensionHandler.onTenantUnSubscribedEvent(tenantUnSubscribedEvent); - } catch (Exception e) { - if (log.isErrorEnabled()) { - log.error("Error processing tenant unSubscribed event", e); - } - } finally { - TenantManager.releaseReadLock(); - } - } - }); - - tenantEventReceiver.execute(); - if (log.isInfoEnabled()) { - log.info("Tenant event message receiver thread started"); + + if (log.isDebugEnabled()) { + log.debug("skipping registerTenantEventListeners before X"); + } + + eventListenerns.startTenantEventReceiver(); + + if (log.isDebugEnabled()) { + log.debug("registerTenantEventListeners after"); } } protected void validateRequiredSystemProperties() { + log.debug("MArtin: validating system properties"); String jndiPropertiesDir = System.getProperty(CartridgeAgentConstants.JNDI_PROPERTIES_DIR); + + log.debug("MArtin: validating system properties 0b: " + jndiPropertiesDir); if (StringUtils.isBlank(jndiPropertiesDir)) { if (log.isErrorEnabled()) { log.error(String.format("System property not found: %s", CartridgeAgentConstants.JNDI_PROPERTIES_DIR)); } return; } + + log.debug("MArtin: validating system properties 1"); String payloadPath = System.getProperty(CartridgeAgentConstants.PARAM_FILE_PATH); + log.debug("MArtin: validating system properties 1b: " + payloadPath); if (StringUtils.isBlank(payloadPath)) { if (log.isErrorEnabled()) { log.error(String.format("System property not found: %s", CartridgeAgentConstants.PARAM_FILE_PATH)); } return; } + + log.debug("MArtin: validating system properties 2"); String extensionsDir = System.getProperty(CartridgeAgentConstants.EXTENSIONS_DIR); + log.debug("MArtin: validating system properties 2b: " + extensionsDir); if (StringUtils.isBlank(extensionsDir)) { if (log.isWarnEnabled()) { log.warn(String.format("System property not found: %s", CartridgeAgentConstants.EXTENSIONS_DIR)); } } + + + log.debug("MArtin: validating system properties 3"); } private static void publishLogs(LogPublisherManager logPublisherManager) { http://git-wip-us.apache.org/repos/asf/stratos/blob/8654829f/components/org.apache.stratos.cartridge.agent/src/main/java/org/apache/stratos/cartridge/agent/CartridgeAgentEventListeners.java ---------------------------------------------------------------------- diff --git a/components/org.apache.stratos.cartridge.agent/src/main/java/org/apache/stratos/cartridge/agent/CartridgeAgentEventListeners.java b/components/org.apache.stratos.cartridge.agent/src/main/java/org/apache/stratos/cartridge/agent/CartridgeAgentEventListeners.java new file mode 100644 index 0000000..b64e32d --- /dev/null +++ b/components/org.apache.stratos.cartridge.agent/src/main/java/org/apache/stratos/cartridge/agent/CartridgeAgentEventListeners.java @@ -0,0 +1,418 @@ +package org.apache.stratos.cartridge.agent; + +import org.apache.commons.logging.LogFactory; +import org.apache.commons.logging.Log; +import org.apache.stratos.cartridge.agent.config.CartridgeAgentConfiguration; +import org.apache.stratos.cartridge.agent.extensions.DefaultExtensionHandler; +import org.apache.stratos.cartridge.agent.extensions.ExtensionHandler; +import org.apache.stratos.common.threading.StratosThreadPool; +import org.apache.stratos.messaging.event.Event; +import org.apache.stratos.messaging.event.instance.notifier.ArtifactUpdatedEvent; +import org.apache.stratos.messaging.event.instance.notifier.InstanceCleanupClusterEvent; +import org.apache.stratos.messaging.event.instance.notifier.InstanceCleanupMemberEvent; +import org.apache.stratos.messaging.event.tenant.CompleteTenantEvent; +import org.apache.stratos.messaging.event.tenant.TenantSubscribedEvent; +import org.apache.stratos.messaging.event.tenant.TenantUnSubscribedEvent; +import org.apache.stratos.messaging.event.topology.*; +import org.apache.stratos.messaging.listener.instance.notifier.ArtifactUpdateEventListener; +import org.apache.stratos.messaging.listener.instance.notifier.InstanceCleanupClusterEventListener; +import org.apache.stratos.messaging.listener.instance.notifier.InstanceCleanupMemberEventListener; +import org.apache.stratos.messaging.listener.tenant.CompleteTenantEventListener; +import org.apache.stratos.messaging.listener.tenant.TenantSubscribedEventListener; +import org.apache.stratos.messaging.listener.tenant.TenantUnSubscribedEventListener; +import org.apache.stratos.messaging.listener.topology.*; +import org.apache.stratos.messaging.message.receiver.instance.notifier.InstanceNotifierEventReceiver; +import org.apache.stratos.messaging.message.receiver.tenant.TenantEventReceiver; +import org.apache.stratos.messaging.message.receiver.tenant.TenantManager; +import org.apache.stratos.messaging.message.receiver.topology.TopologyEventReceiver; +import org.apache.stratos.messaging.message.receiver.topology.TopologyManager; + +import java.util.concurrent.ExecutorService; + + +public class CartridgeAgentEventListeners +{ + + private static final Log log = LogFactory.getLog(CartridgeAgentEventListeners.class); + + private InstanceNotifierEventReceiver instanceNotifierEventReceiver; + private TopologyEventReceiver topologyEventReceiver; + private TenantEventReceiver tenantEventReceiver; + + private ExtensionHandler extensionHandler; + private static final ExecutorService instanceNotifierExecutorService = + StratosThreadPool.getExecutorService("CARTRIDGEAGENT_NOTIFIER_EVENT_LISTENERS", 1); + + private static final ExecutorService topologyEventExecutorService = + StratosThreadPool.getExecutorService("CARTRIDGEAGENT_TOPOLOGY_EVENT_LISTENERS", 1); + + private static final ExecutorService tenantEventExecutorService = + StratosThreadPool.getExecutorService("CARTRIDGEAGENT_TENANT_EVENT_LISTENERS", 1); + + public CartridgeAgentEventListeners() { + if (log.isDebugEnabled()) { + log.debug("creating cartridgeAgentEventListeners ... "); + } + this.topologyEventReceiver = new TopologyEventReceiver(); + this.topologyEventReceiver.setExecutorService(topologyEventExecutorService); + + this.instanceNotifierEventReceiver = new InstanceNotifierEventReceiver(); + + this.tenantEventReceiver = new TenantEventReceiver(); + this.tenantEventReceiver.setExecutorService(tenantEventExecutorService); + + extensionHandler = new DefaultExtensionHandler(); + + addInstanceNotifierEventListeners(); + addTopologyEventListeners(); + addTenantEventListeners(); + + if (log.isDebugEnabled()) { + log.debug("creating cartridgeAgentEventListeners ... done "); + } + + } + + public void startTopologyEventReceiver() { + + if (log.isDebugEnabled()) { + log.debug("Starting cartridge agent topology event message receiver"); + } + + topologyEventExecutorService.submit(new Runnable() { + @Override + public void run() { + topologyEventReceiver.execute(); + } + }); + + if (log.isInfoEnabled()) + { + log.info("Cartridge agent topology receiver thread started, waiting for event messages ..."); + } + + } + + public void startInstanceNotifierReceiver() { + + if (log.isDebugEnabled()) { + log.debug("Starting cartridge agent instance notifier event message receiver"); + } + + instanceNotifierExecutorService.submit(new Runnable() { + @Override + public void run() { + instanceNotifierEventReceiver.execute(); + } + }); + + if (log.isDebugEnabled()) { + log.debug("Cartridge agent Instance notifier event message receiver started, waiting for event messages ..."); + } + } + + public void startTenantEventReceiver() { + + if (log.isDebugEnabled()) { + log.debug("Starting cartridge agent tenant event message receiver"); + } + + topologyEventExecutorService.submit(new Runnable() { + @Override + public void run() { + topologyEventReceiver.execute(); + } + }); + + if (log.isInfoEnabled()) { + log.info("Cartridge agent tenant receiver thread started, waiting for event messages ..."); + } + + } + + + private void addInstanceNotifierEventListeners() { + instanceNotifierEventReceiver.addEventListener(new ArtifactUpdateEventListener() { + @Override + protected void onEvent(Event event) { + try { + extensionHandler.onArtifactUpdatedEvent((ArtifactUpdatedEvent) event); + } catch (Exception e) { + if (log.isErrorEnabled()) { + log.error("Error processing artifact update event", e); + } + } + } + }); + + instanceNotifierEventReceiver.addEventListener(new InstanceCleanupMemberEventListener() { + @Override + protected void onEvent(Event event) { + try { + String memberIdInPayload = CartridgeAgentConfiguration.getInstance().getMemberId(); + InstanceCleanupMemberEvent instanceCleanupMemberEvent = (InstanceCleanupMemberEvent) event; + if (memberIdInPayload.equals(instanceCleanupMemberEvent.getMemberId())) { + extensionHandler.onInstanceCleanupMemberEvent(instanceCleanupMemberEvent); + } + } catch (Exception e) { + if (log.isErrorEnabled()) { + log.error("Error processing instance cleanup member event", e); + } + } + + } + }); + + instanceNotifierEventReceiver.addEventListener(new InstanceCleanupClusterEventListener() { + @Override + protected void onEvent(Event event) { + String clusterIdInPayload = CartridgeAgentConfiguration.getInstance().getClusterId(); + InstanceCleanupClusterEvent instanceCleanupClusterEvent = (InstanceCleanupClusterEvent) event; + if (clusterIdInPayload.equals(instanceCleanupClusterEvent.getClusterId())) { + extensionHandler.onInstanceCleanupClusterEvent(instanceCleanupClusterEvent); + } + } + }); + + + if(log.isInfoEnabled()) { + log.info("Instance notifier event listener added ... "); + } + } + + private void addTopologyEventListeners() { + topologyEventReceiver.addEventListener(new MemberCreatedEventListener() { + @Override + protected void onEvent(Event event) { + try { + boolean initialized = CartridgeAgentConfiguration.getInstance().isInitialized(); + if (initialized) { + // no need to process this event, if the member is initialized. + return; + } + TopologyManager.acquireReadLock(); + if (log.isDebugEnabled()) { + log.debug("Member created event received"); + } + MemberCreatedEvent memberCreatedEvent = (MemberCreatedEvent) event; + extensionHandler.onMemberCreatedEvent(memberCreatedEvent); + } catch (Exception e) { + if (log.isErrorEnabled()) { + log.error("Error processing member created event", e); + } + } finally { + TopologyManager.releaseReadLock(); + } + } + }); + + topologyEventReceiver.addEventListener(new MemberActivatedEventListener() { + @Override + protected void onEvent(Event event) { + boolean initialized = CartridgeAgentConfiguration.getInstance().isInitialized(); + if (!initialized) { + return; + } + try { + TopologyManager.acquireReadLock(); + if (log.isDebugEnabled()) { + log.debug("Member activated event received"); + } + MemberActivatedEvent memberActivatedEvent = (MemberActivatedEvent) event; + extensionHandler.onMemberActivatedEvent(memberActivatedEvent); + } catch (Exception e) { + if (log.isErrorEnabled()) { + log.error("Error processing member activated event", e); + } + } finally { + TopologyManager.releaseReadLock(); + } + } + }); + + topologyEventReceiver.addEventListener(new MemberTerminatedEventListener() { + @Override + protected void onEvent(Event event) { + boolean initialized = CartridgeAgentConfiguration.getInstance().isInitialized(); + if (!initialized) { + return; + } + try { + TopologyManager.acquireReadLock(); + if (log.isDebugEnabled()) { + log.debug("Member terminated event received"); + } + MemberTerminatedEvent memberTerminatedEvent = (MemberTerminatedEvent) event; + extensionHandler.onMemberTerminatedEvent(memberTerminatedEvent); + } catch (Exception e) { + if (log.isErrorEnabled()) { + log.error("Error processing member terminated event", e); + } + } finally { + TopologyManager.releaseReadLock(); + } + } + }); + + topologyEventReceiver.addEventListener(new MemberSuspendedEventListener() { + @Override + protected void onEvent(Event event) { + boolean initialized = CartridgeAgentConfiguration.getInstance().isInitialized(); + if (!initialized) { + return; + } + try { + TopologyManager.acquireReadLock(); + if (log.isDebugEnabled()) { + log.debug("Member suspended event received"); + } + MemberSuspendedEvent memberSuspendedEvent = (MemberSuspendedEvent) event; + extensionHandler.onMemberSuspendedEvent(memberSuspendedEvent); + } catch (Exception e) { + if (log.isErrorEnabled()) { + log.error("Error processing member suspended event", e); + } + } finally { + TopologyManager.releaseReadLock(); + } + } + }); + + topologyEventReceiver.addEventListener(new CompleteTopologyEventListener() { + @Override + protected void onEvent(Event event) { + boolean initialized = CartridgeAgentConfiguration.getInstance().isInitialized(); + if (!initialized) { + try { + TopologyManager.acquireReadLock(); + if (log.isDebugEnabled()) { + log.debug("Complete topology event received"); + } + CompleteTopologyEvent completeTopologyEvent = (CompleteTopologyEvent) event; + extensionHandler.onCompleteTopologyEvent(completeTopologyEvent); + } catch (Exception e) { + if (log.isErrorEnabled()) { + log.error("Error processing complete topology event", e); + } + } finally { + TopologyManager.releaseReadLock(); + } + } + } + }); + + topologyEventReceiver.addEventListener(new MemberStartedEventListener() { + @Override + protected void onEvent(Event event) { + boolean initialized = CartridgeAgentConfiguration.getInstance().isInitialized(); + if (!initialized) { + return; + } + try { + TopologyManager.acquireReadLock(); + if (log.isDebugEnabled()) { + log.debug("Member started event received"); + } + MemberStartedEvent memberStartedEvent = (MemberStartedEvent) event; + extensionHandler.onMemberStartedEvent(memberStartedEvent); + } catch (Exception e) { + if (log.isErrorEnabled()) { + log.error("Error processing member started event", e); + } + } finally { + TopologyManager.releaseReadLock(); + } + } + }); + + if(log.isInfoEnabled()) { + log.info("Topology event listener added ... "); + } + } + + private void addTenantEventListeners() { + + tenantEventReceiver.addEventListener(new CompleteTenantEventListener() { + private boolean initialized; + @Override + protected void onEvent(Event event) { + if (!initialized) { + try { + TenantManager.acquireReadLock(); + if (log.isDebugEnabled()) { + log.debug("Complete tenant event received"); + } + CompleteTenantEvent completeTenantEvent = (CompleteTenantEvent) event; + extensionHandler.onCompleteTenantEvent(completeTenantEvent); + initialized = true; + } catch (Exception e) { + if (log.isErrorEnabled()) { + log.error("Error processing complete tenant event", e); + } + } finally { + TenantManager.releaseReadLock(); + } + + } else { + if (log.isInfoEnabled()) { + log.info("Complete tenant event updating task disabled"); + } + } + } + }); + + tenantEventReceiver.addEventListener(new TenantSubscribedEventListener() { + @Override + protected void onEvent(Event event) { + try { + TenantManager.acquireReadLock(); + if (log.isDebugEnabled()) { + log.debug("Tenant subscribed event received"); + } + TenantSubscribedEvent tenantSubscribedEvent = (TenantSubscribedEvent) event; + extensionHandler.onTenantSubscribedEvent(tenantSubscribedEvent); + } catch (Exception e) { + if (log.isErrorEnabled()) { + log.error("Error processing tenant subscribed event", e); + } + } finally { + TenantManager.releaseReadLock(); + } + } + }); + + tenantEventReceiver.addEventListener(new TenantUnSubscribedEventListener() { + @Override + protected void onEvent(Event event) { + try { + TenantManager.acquireReadLock(); + if (log.isDebugEnabled()) { + log.debug("Tenant unSubscribed event received"); + } + TenantUnSubscribedEvent tenantUnSubscribedEvent = (TenantUnSubscribedEvent) event; + extensionHandler.onTenantUnSubscribedEvent(tenantUnSubscribedEvent); + } catch (Exception e) { + if (log.isErrorEnabled()) { + log.error("Error processing tenant unSubscribed event", e); + } + } finally { + TenantManager.releaseReadLock(); + } + } + }); + + if(log.isInfoEnabled()) { + log.info("Tenant event listener added ... "); + } + } + + /** + * Terminate load balancer topology receiver thread. + */ + + public void terminate() { + topologyEventReceiver.terminate(); + } + +} + http://git-wip-us.apache.org/repos/asf/stratos/blob/8654829f/components/org.apache.stratos.cartridge.agent/src/main/java/org/apache/stratos/cartridge/agent/config/CartridgeAgentConfiguration.java ---------------------------------------------------------------------- diff --git a/components/org.apache.stratos.cartridge.agent/src/main/java/org/apache/stratos/cartridge/agent/config/CartridgeAgentConfiguration.java b/components/org.apache.stratos.cartridge.agent/src/main/java/org/apache/stratos/cartridge/agent/config/CartridgeAgentConfiguration.java index 1ae1a79..70f62ef 100644 --- a/components/org.apache.stratos.cartridge.agent/src/main/java/org/apache/stratos/cartridge/agent/config/CartridgeAgentConfiguration.java +++ b/components/org.apache.stratos.cartridge.agent/src/main/java/org/apache/stratos/cartridge/agent/config/CartridgeAgentConfiguration.java @@ -74,22 +74,34 @@ public class CartridgeAgentConfiguration { private String instanceId; private String clusterInstanceId; private String applicationId; + private String dependantClusterId; private CartridgeAgentConfiguration() { parameters = loadParametersFile(); try { + //self.application_id = self.read_property(cartridgeagentconstants.APPLICATION_ID) + applicationId = readApplicationId(); serviceGroup = readServiceGroup(); isClustered = readClustering(); + //self.service_name = self.read_property(cartridgeagentconstants.SERVICE_NAME) serviceName = readParameterValue(CartridgeAgentConstants.SERVICE_NAME); + //self.cluster_id = self.read_property(cartridgeagentconstants.CLUSTER_ID) clusterId = readParameterValue(CartridgeAgentConstants.CLUSTER_ID); + //self.network_partition_id = self.read_property(cartridgeagentconstants.NETWORK_PARTITION_ID, False) networkPartitionId = readParameterValue(CartridgeAgentConstants.NETWORK_PARTITION_ID); + //self.partition_id = self.read_property(cartridgeagentconstants.PARTITION_ID, False) partitionId = readParameterValue(CartridgeAgentConstants.PARTITION_ID); + //self.member_id = self.read_property(cartridgeagentconstants.MEMBER_ID) memberId = readMemberIdValue(CartridgeAgentConstants.MEMBER_ID); + //self.cartridge_key = self.read_property(cartridgeagentconstants.CARTRIDGE_KEY) cartridgeKey = readParameterValue(CartridgeAgentConstants.CARTRIDGE_KEY); + //self.app_path = self.read_property(cartridgeagentconstants.APPLICATION_PATH, False) appPath = readParameterValue(CartridgeAgentConstants.APP_PATH); + //self.repo_url = self.read_property(cartridgeagentconstants.REPO_URL, False) repoUrl = readParameterValue(CartridgeAgentConstants.REPO_URL); + //self.ports = str(self.read_property(cartridgeagentconstants.PORTS)).split("|") ports = readPorts(); logFilePaths = readLogFilePaths(); isMultitenant = readMultitenant(CartridgeAgentConstants.MULTITENANT); @@ -114,6 +126,13 @@ public class CartridgeAgentConfiguration { isPrimary = readIsPrimary(); kubernetesClusterId = readKubernetesClusterIdValue(CartridgeAgentConstants.KUBERNETES_CLUSTER_ID); + //self.cluster_instance_id= self.read_property(cartridgeagentconstants.CLUSTER_INSTANCE_ID) + clusterInstanceId = readClusterInstanceId(); + //self.dependant_cluster_id = self.read_property(cartridgeagentconstants.DEPENDENCY_CLUSTER_IDS, False) + dependantClusterId = readDependantClusterId(); + + + } catch (ParameterNotFoundException e) { throw new RuntimeException(e); } @@ -123,6 +142,9 @@ public class CartridgeAgentConfiguration { } if (log.isDebugEnabled()) { + log.debug(String.format("application-id: %s", applicationId)); + log.debug(String.format("service-group: %s", serviceGroup)); + log.debug(String.format("cluster-id: %s", clusterId)); log.debug(String.format("service-name: %s", serviceName)); log.debug(String.format("cluster-id: %s", clusterId)); log.debug(String.format("network-partition-id: %s", networkPartitionId)); @@ -134,6 +156,8 @@ public class CartridgeAgentConfiguration { log.debug(String.format("ports: %s", ports.toString())); log.debug(String.format("lb-private-ip: %s", lbPrivateIp)); log.debug(String.format("lb-public-ip: %s", lbPublicIp)); + log.debug(String.format("cluster-instance-id: %s", clusterInstanceId)); + log.debug(String.format("dependant-cluster-id: %s", dependantClusterId)); } } @@ -364,7 +388,15 @@ public class CartridgeAgentConfiguration { return parameters; } - + + private String readApplicationId() { + if (parameters.containsKey(CartridgeAgentConstants.APPLICATION_ID)) { + return parameters.get(CartridgeAgentConstants.APPLICATION_ID); + } else { + return null; + } + } + private String readServiceGroup() { if (parameters.containsKey(CartridgeAgentConstants.SERVICE_GROUP)) { return parameters.get(CartridgeAgentConstants.SERVICE_GROUP); @@ -372,6 +404,22 @@ public class CartridgeAgentConfiguration { return null; } } + + private String readClusterInstanceId() { + if (parameters.containsKey(CartridgeAgentConstants.CLUSTER_INSTANCE_ID)) { + return parameters.get(CartridgeAgentConstants.CLUSTER_INSTANCE_ID); + } else { + return null; + } + } + + private String readDependantClusterId() { + if (parameters.containsKey(CartridgeAgentConstants.DEPENDENCY_CLUSTER_IDS)) { + return parameters.get(CartridgeAgentConstants.DEPENDENCY_CLUSTER_IDS); + } else { + return null; + } + } private String readClustering() { if (parameters.containsKey(CartridgeAgentConstants.CLUSTERING)) { http://git-wip-us.apache.org/repos/asf/stratos/blob/8654829f/components/org.apache.stratos.cartridge.agent/src/main/java/org/apache/stratos/cartridge/agent/event/publisher/CartridgeAgentEventPublisher.java ---------------------------------------------------------------------- diff --git a/components/org.apache.stratos.cartridge.agent/src/main/java/org/apache/stratos/cartridge/agent/event/publisher/CartridgeAgentEventPublisher.java b/components/org.apache.stratos.cartridge.agent/src/main/java/org/apache/stratos/cartridge/agent/event/publisher/CartridgeAgentEventPublisher.java index 22f4ca2..cfdcd16 100644 --- a/components/org.apache.stratos.cartridge.agent/src/main/java/org/apache/stratos/cartridge/agent/event/publisher/CartridgeAgentEventPublisher.java +++ b/components/org.apache.stratos.cartridge.agent/src/main/java/org/apache/stratos/cartridge/agent/event/publisher/CartridgeAgentEventPublisher.java @@ -49,13 +49,40 @@ public class CartridgeAgentEventPublisher { log.info("Publishing instance started event"); } InstanceStartedEvent event = new InstanceStartedEvent( + //application_id = CartridgeAgentConfiguration().application_id CartridgeAgentConfiguration.getInstance().getApplicationId(), + //service_name = CartridgeAgentConfiguration().service_name CartridgeAgentConfiguration.getInstance().getServiceName(), + //cluster_id = CartridgeAgentConfiguration().cluster_id CartridgeAgentConfiguration.getInstance().getClusterId(), + // member_id = CartridgeAgentConfiguration().member_id CartridgeAgentConfiguration.getInstance().getMemberId(), + + //cluster_instance_id = CartridgeAgentConfiguration().cluster_instance_id CartridgeAgentConfiguration.getInstance().getClusterInstanceId(), + //network_partition_id = CartridgeAgentConfiguration().network_partition_id CartridgeAgentConfiguration.getInstance().getNetworkPartitionId(), + //partition_id = CartridgeAgentConfiguration().partition_id CartridgeAgentConfiguration.getInstance().getPartitionId()); + + /* + * + + public InstanceStartedEvent(String applicationId, String serviceName, String clusterId, String memberId, + String clusterInstanceId, String networkPartitionId, String partitionId) + + //instance_id = CartridgeAgentConfiguration().instance_id + + + + + instance_started_event = InstanceStartedEvent(application_id, service_name, cluster_id, cluster_instance_id, member_id, + instance_id, network_partition_id, partition_id) + publisher = get_publisher(cartridgeagentconstants.INSTANCE_STATUS_TOPIC + cartridgeagentconstants.INSTANCE_STARTED_EVENT) + publisher.publish(instance_started_event) + started = True + log.info("Instance started event published") + */ String topic = Util.getMessageTopicName(event); EventPublisher eventPublisher = EventPublisherPool http://git-wip-us.apache.org/repos/asf/stratos/blob/8654829f/components/org.apache.stratos.cartridge.agent/src/main/java/org/apache/stratos/cartridge/agent/util/CartridgeAgentConstants.java ---------------------------------------------------------------------- diff --git a/components/org.apache.stratos.cartridge.agent/src/main/java/org/apache/stratos/cartridge/agent/util/CartridgeAgentConstants.java b/components/org.apache.stratos.cartridge.agent/src/main/java/org/apache/stratos/cartridge/agent/util/CartridgeAgentConstants.java index 995136b..b1db135 100644 --- a/components/org.apache.stratos.cartridge.agent/src/main/java/org/apache/stratos/cartridge/agent/util/CartridgeAgentConstants.java +++ b/components/org.apache.stratos.cartridge.agent/src/main/java/org/apache/stratos/cartridge/agent/util/CartridgeAgentConstants.java @@ -38,6 +38,7 @@ public class CartridgeAgentConstants implements Serializable{ public static final String CARTRIDGE_KEY = "CARTRIDGE_KEY"; public static final String APP_PATH = "APP_PATH"; + public static final String APPLICATION_ID = "APPLICATION_ID"; public static final String SERVICE_GROUP = "SERIVCE_GROUP"; public static final String SERVICE_NAME = "SERVICE_NAME"; public static final String CLUSTER_ID = "CLUSTER_ID"; @@ -51,6 +52,8 @@ public class CartridgeAgentConstants implements Serializable{ public static final String DEPLOYMENT = "DEPLOYMENT"; public static final String MANAGER_SERVICE_TYPE = "MANAGER_SERVICE_TYPE"; public static final String WORKER_SERVICE_TYPE = "WORKER_SERVICE_TYPE"; + public static final String DEPENDENCY_CLUSTER_IDS = "DEPENDENCY_CLUSTER_IDS"; + public static final String CLUSTER_INSTANCE_ID = "CLUSTER_INSTANCE_ID"; // stratos.sh environment variables keys public static final String LOG_FILE_PATHS = "LOG_FILE_PATHS";
