Repository: karaf-cellar Updated Branches: refs/heads/cellar-3.0.x f71624e5c -> 7c3e1c703
[KARAF-3981] Improve features synchronizer Project: http://git-wip-us.apache.org/repos/asf/karaf-cellar/repo Commit: http://git-wip-us.apache.org/repos/asf/karaf-cellar/commit/7c3e1c70 Tree: http://git-wip-us.apache.org/repos/asf/karaf-cellar/tree/7c3e1c70 Diff: http://git-wip-us.apache.org/repos/asf/karaf-cellar/diff/7c3e1c70 Branch: refs/heads/cellar-3.0.x Commit: 7c3e1c703b4898f7f42181cf0d9d90b05f07fa4e Parents: f71624e Author: Jean-Baptiste Onofré <[email protected]> Authored: Sun Sep 13 07:52:52 2015 +0200 Committer: Jean-Baptiste Onofré <[email protected]> Committed: Mon Sep 14 08:40:09 2015 +0200 ---------------------------------------------------------------------- assembly/src/main/resources/groups.cfg | 18 +-- .../cellar/features/FeaturesEventHandler.java | 6 + .../cellar/features/FeaturesSynchronizer.java | 122 +++++++++++++++---- .../cellar/features/LocalFeaturesListener.java | 6 +- .../resources/OSGI-INF/blueprint/blueprint.xml | 1 + 5 files changed, 118 insertions(+), 35 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/karaf-cellar/blob/7c3e1c70/assembly/src/main/resources/groups.cfg ---------------------------------------------------------------------- diff --git a/assembly/src/main/resources/groups.cfg b/assembly/src/main/resources/groups.cfg index 62a0cfe..bd64cd6 100644 --- a/assembly/src/main/resources/groups.cfg +++ b/assembly/src/main/resources/groups.cfg @@ -17,17 +17,13 @@ default.bundle.blacklist.outbound = *.xml default.config.whitelist.inbound = * default.config.whitelist.outbound = * default.config.blacklist.inbound = org.apache.felix.fileinstall*, \ - org.apache.karaf.cellar*, \ org.apache.karaf.management, \ org.apache.karaf.shell, \ - org.ops4j.pax.logging, \ org.ops4j.pax.web, \ org.apache.aries.transaction default.config.blacklist.outbound = org.apache.felix.fileinstall*, \ - org.apache.karaf.cellar*, \ org.apache.karaf.management, \ org.apache.karaf.shell, \ - org.ops4j.pax.logging, \ org.ops4j.pax.web, \ org.apache.aries.transaction @@ -43,11 +39,15 @@ default.feature.blacklist.outbound = none # The following properties define the behavior to use when the node joins the cluster (the usage of the bootstrap # synchronizer), per cluster group and per resource. # The following values are accepted: -# disabled: means that the synchronizer is not used, meaning the node or the cluster are not updated at all -# cluster: if the node is the first one in the cluster, it pushes its local state to the cluster, else it's not the -# first node of the cluster, the node will update its local state with the cluster one (meaning that the cluster -# is the master) -# node: in this case, the node is the master, it means that the cluster state will be overwritten by the node state. +# disabled: means that the synchronizer doesn't sync cluster group and node states +# cluster: the synchronizer retrieves the state from the cluster group first (pull first), and push the node the state +# to the cluster group after (push after) +# node: the synchronizer push the node state to the cluster group (push first), and pull the state from the cluster group + after (pull after) +# clusterOnly: the cluster is the "master", the node only retrieves and applies the cluster group state, nothing is +# pushed to the cluster group +# nodeOnly: the node is the "master", the node pushes his state to the cluster group, nothing is pulled from the +# cluster group # default.bundle.sync = cluster default.config.sync = cluster http://git-wip-us.apache.org/repos/asf/karaf-cellar/blob/7c3e1c70/features/src/main/java/org/apache/karaf/cellar/features/FeaturesEventHandler.java ---------------------------------------------------------------------- diff --git a/features/src/main/java/org/apache/karaf/cellar/features/FeaturesEventHandler.java b/features/src/main/java/org/apache/karaf/cellar/features/FeaturesEventHandler.java index eb706fa..cab9fc7 100644 --- a/features/src/main/java/org/apache/karaf/cellar/features/FeaturesEventHandler.java +++ b/features/src/main/java/org/apache/karaf/cellar/features/FeaturesEventHandler.java @@ -72,6 +72,12 @@ public class FeaturesEventHandler extends FeaturesSupport implements EventHandle return; } + // check if it's not a "local" event + if (event.getSourceNode() != null && event.getSourceNode().getId().equalsIgnoreCase(clusterManager.getNode().getId())) { + LOGGER.trace("CELLAR FEATURE: cluster event is local (coming from local synchronizer or listener)"); + return; + } + String name = event.getName(); String version = event.getVersion(); if (isAllowed(event.getSourceGroup(), Constants.CATEGORY, name, EventType.INBOUND) || event.getForce()) { http://git-wip-us.apache.org/repos/asf/karaf-cellar/blob/7c3e1c70/features/src/main/java/org/apache/karaf/cellar/features/FeaturesSynchronizer.java ---------------------------------------------------------------------- diff --git a/features/src/main/java/org/apache/karaf/cellar/features/FeaturesSynchronizer.java b/features/src/main/java/org/apache/karaf/cellar/features/FeaturesSynchronizer.java index 7ab60b7..8d95b57 100644 --- a/features/src/main/java/org/apache/karaf/cellar/features/FeaturesSynchronizer.java +++ b/features/src/main/java/org/apache/karaf/cellar/features/FeaturesSynchronizer.java @@ -16,9 +16,13 @@ package org.apache.karaf.cellar.features; import org.apache.karaf.cellar.core.Configurations; import org.apache.karaf.cellar.core.Group; import org.apache.karaf.cellar.core.Synchronizer; +import org.apache.karaf.cellar.core.control.SwitchStatus; +import org.apache.karaf.cellar.core.event.EventProducer; import org.apache.karaf.cellar.core.event.EventType; import org.apache.karaf.features.Feature; +import org.apache.karaf.features.FeatureEvent; import org.apache.karaf.features.Repository; +import org.apache.karaf.features.RepositoryEvent; import org.osgi.service.cm.Configuration; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -37,6 +41,12 @@ public class FeaturesSynchronizer extends FeaturesSupport implements Synchronize private static final transient Logger LOGGER = LoggerFactory.getLogger(FeaturesSynchronizer.class); + private EventProducer eventProducer; + + public void setEventProducer(EventProducer eventProducer) { + this.eventProducer = eventProducer; + } + @Override public void init() { Set<Group> groups = groupManager.listLocalGroups(); @@ -60,19 +70,32 @@ public class FeaturesSynchronizer extends FeaturesSupport implements Synchronize @Override public void sync(Group group) { String policy = getSyncPolicy(group); - if (policy != null && policy.equalsIgnoreCase("cluster")) { - LOGGER.debug("CELLAR FEATURE: sync policy is set as 'cluster' for cluster group " + group.getName()); - if (clusterManager.listNodesByGroup(group).size() == 1 && clusterManager.listNodesByGroup(group).contains(clusterManager.getNode())) { - LOGGER.debug("CELLAR FEATURE: node is the first and only member of the group, pushing state"); - push(group); - } else { - LOGGER.debug("CELLAR FEATURE: pulling state"); - pull(group); - } + if (policy == null) { + LOGGER.warn("CELLAR FEATURE: sync policy is not defined for cluster group {}", group.getName()); } - if (policy != null && policy.equalsIgnoreCase("node")) { - LOGGER.debug("CELLAR FEATURE: sync policy is set as 'node' for cluster group " + group.getName()); + if (policy.equalsIgnoreCase("cluster")) { + LOGGER.debug("CELLAR FEATURE: sync policy set as 'cluster' for cluster group {}", group.getName()); + LOGGER.debug("CELLAR FEATURE: updating node from the cluster (pull first)"); + pull(group); + LOGGER.debug("CELLAR FEATURE: updating cluster from the local node (push after)"); + push(group); + } else if (policy.equalsIgnoreCase("node")) { + LOGGER.debug("CELLAR FEATURE: sync policy set as 'node' for cluster group {}", group.getName()); + LOGGER.debug("CELLAR FEATURE: updating cluster from the local node (push first)"); push(group); + LOGGER.debug("CELLAR FEATURE: updating node from the cluster (pull after)"); + pull(group); + } else if (policy.equalsIgnoreCase("clusterOnly")) { + LOGGER.debug("CELLAR FEATURE: sync policy set as 'clusterOnly' for cluster group " + group.getName()); + LOGGER.debug("CELLAR FEATURE: updating node from the cluster (pull only)"); + pull(group); + } else if (policy.equalsIgnoreCase("nodeOnly")) { + LOGGER.debug("CELLAR FEATURE: sync policy set as 'nodeOnly' for cluster group " + group.getName()); + LOGGER.debug("CELLAR FEATURE: updating cluster from the local node (push only)"); + push(group); + } else { + LOGGER.debug("CELLAR FEATURE: sync policy set as 'disabled' for cluster group " + group.getName()); + LOGGER.debug("CELLAR FEATURE: no sync"); } } @@ -100,7 +123,7 @@ public class FeaturesSynchronizer extends FeaturesSupport implements Synchronize if (!isRepositoryRegisteredLocally(url)) { LOGGER.debug("CELLAR FEATURE: adding repository {}", url); featuresService.addRepository(new URI(url)); - } + } // TODO uninstall local features repositories not on the cluster ? } catch (MalformedURLException e) { LOGGER.error("CELLAR FEATURE: failed to add repository URL {} (malformed)", url, e); } catch (Exception e) { @@ -134,7 +157,7 @@ public class FeaturesSynchronizer extends FeaturesSupport implements Synchronize } catch (Exception e) { LOGGER.error("CELLAR FEATURE: failed to install feature {}/{} ", new Object[]{state.getName(), state.getVersion()}, e); } - } + } // TODO uninstall local features not on the cluster ? } else LOGGER.trace("CELLAR FEATURE: feature {} is marked BLOCKED INBOUND for cluster group {}", name, groupName); } } @@ -151,6 +174,12 @@ public class FeaturesSynchronizer extends FeaturesSupport implements Synchronize */ @Override public void push(Group group) { + + if (eventProducer.getSwitch().getStatus().equals(SwitchStatus.OFF)) { + LOGGER.warn("CELLAR FEATURE: cluster event producer is OFF"); + return; + } + if (group != null) { String groupName = group.getName(); LOGGER.debug("CELLAR FEATURE: pushing features repositories and features in cluster group {}", groupName); @@ -175,11 +204,21 @@ public class FeaturesSynchronizer extends FeaturesSupport implements Synchronize // push features repositories to the cluster group if (repositoryList != null && repositoryList.length > 0) { for (Repository repository : repositoryList) { - if (!clusterRepositories.containsKey(repository.getURI().toString())) { - clusterRepositories.put(repository.getURI().toString(), repository.getName()); - LOGGER.debug("CELLAR FEATURE: pushing repository {} in cluster group {}", repository.getName(), groupName); - } else { - LOGGER.debug("CELLAR FEATURE: repository {} is already in cluster group {}", repository.getName(), groupName); + try { + if (!clusterRepositories.containsKey(repository.getURI().toString())) { + LOGGER.debug("CELLAR FEATURE: pushing repository {} in cluster group {}", repository.getName(), groupName); + // updating cluster state + clusterRepositories.put(repository.getURI().toString(), repository.getName()); + // sending cluster event + ClusterRepositoryEvent event = new ClusterRepositoryEvent(repository.getURI().toString(), RepositoryEvent.EventType.RepositoryAdded); + event.setSourceGroup(group); + event.setSourceNode(clusterManager.getNode()); + eventProducer.produce(event); + } else { + LOGGER.debug("CELLAR FEATURE: repository {} is already in cluster group {}", repository.getName(), groupName); + } + } catch (Exception e) { + LOGGER.warn("CELLAR FEATURE: can't add repository", e); } } } @@ -188,12 +227,47 @@ public class FeaturesSynchronizer extends FeaturesSupport implements Synchronize if (featuresList != null && featuresList.length > 0) { for (Feature feature : featuresList) { if (isAllowed(group, Constants.CATEGORY, feature.getName(), EventType.OUTBOUND)) { - FeatureState clusterFeatureState = new FeatureState(); - clusterFeatureState.setName(feature.getName()); - clusterFeatureState.setVersion(feature.getVersion()); - clusterFeatureState.setInstalled(featuresService.isInstalled(feature)); - clusterFeatures.put(feature.getName() + "/" + feature.getVersion(), clusterFeatureState); - LOGGER.debug("CELLAR FEATURE : pushing feature {}/{} to cluster group {}", feature.getName(), feature.getVersion(), groupName); + boolean installed = featuresService.isInstalled(feature); + String key = feature.getName() + "/" + feature.getVersion(); + FeatureState clusterFeature = clusterFeatures.get(key); + if (clusterFeature == null) { + LOGGER.debug("CELLAR FEATURE: adding feature {} to cluster group {}", key, groupName); + // updating cluster state + clusterFeature = new FeatureState(); + clusterFeature.setName(feature.getName()); + clusterFeature.setVersion(feature.getVersion()); + clusterFeature.setInstalled(installed); + clusterFeatures.put(key, clusterFeature); + // sending cluster event + ClusterFeaturesEvent event; + if (installed) { + event = new ClusterFeaturesEvent(feature.getName(), feature.getVersion(), FeatureEvent.EventType.FeatureInstalled); + } else { + event = new ClusterFeaturesEvent(feature.getName(), feature.getVersion(), FeatureEvent.EventType.FeatureUninstalled); + } + event.setSourceGroup(group); + event.setSourceNode(clusterManager.getNode()); + eventProducer.produce(event); + + } else { + if (clusterFeature.getInstalled() != installed) { + // updating cluster state + clusterFeature.setInstalled(installed); + clusterFeatures.put(key, clusterFeature); + // sending cluster event + ClusterFeaturesEvent event; + if (installed) { + event = new ClusterFeaturesEvent(feature.getName(), feature.getVersion(), FeatureEvent.EventType.FeatureInstalled); + } else { + event = new ClusterFeaturesEvent(feature.getName(), feature.getVersion(), FeatureEvent.EventType.FeatureUninstalled); + } + event.setSourceGroup(group); + event.setSourceNode(clusterManager.getNode()); + eventProducer.produce(event); + } else { + LOGGER.debug("CELLAR FEATURE: feature {} already sync on the cluster group {}", key, groupName); + } + } } else { LOGGER.debug("CELLAR FEATURE: feature {} is marked BLOCKED OUTBOUND for cluster group {}", feature.getName(), groupName); } http://git-wip-us.apache.org/repos/asf/karaf-cellar/blob/7c3e1c70/features/src/main/java/org/apache/karaf/cellar/features/LocalFeaturesListener.java ---------------------------------------------------------------------- diff --git a/features/src/main/java/org/apache/karaf/cellar/features/LocalFeaturesListener.java b/features/src/main/java/org/apache/karaf/cellar/features/LocalFeaturesListener.java index 26c5f50..2945c92 100644 --- a/features/src/main/java/org/apache/karaf/cellar/features/LocalFeaturesListener.java +++ b/features/src/main/java/org/apache/karaf/cellar/features/LocalFeaturesListener.java @@ -57,7 +57,7 @@ public class LocalFeaturesListener extends FeaturesSupport implements org.apache public void featureEvent(FeatureEvent event) { if (!isEnabled()) { - LOGGER.debug("CELLAR FEATURE: local listener is disabled"); + LOGGER.trace("CELLAR FEATURE: local listener is disabled"); return; } @@ -95,6 +95,7 @@ public class LocalFeaturesListener extends FeaturesSupport implements org.apache // broadcast the event ClusterFeaturesEvent featureEvent = new ClusterFeaturesEvent(name, version, type); featureEvent.setSourceGroup(group); + featureEvent.setSourceNode(clusterManager.getNode()); eventProducer.produce(featureEvent); } else LOGGER.trace("CELLAR FEATURE: feature {} is marked BLOCKED OUTBOUND for cluster group {}", name, group.getName()); } @@ -111,7 +112,7 @@ public class LocalFeaturesListener extends FeaturesSupport implements org.apache public void repositoryEvent(RepositoryEvent event) { if (!isEnabled()) { - LOGGER.debug("CELLAR FEATURE: local listener is disabled"); + LOGGER.trace("CELLAR FEATURE: local listener is disabled"); return; } @@ -131,6 +132,7 @@ public class LocalFeaturesListener extends FeaturesSupport implements org.apache for (Group group : groups) { ClusterRepositoryEvent clusterRepositoryEvent = new ClusterRepositoryEvent(event.getRepository().getURI().toString(), event.getType()); clusterRepositoryEvent.setSourceGroup(group); + clusterRepositoryEvent.setSourceNode(clusterManager.getNode()); RepositoryEvent.EventType type = event.getType(); Map<String, String> clusterRepositories = clusterManager.getMap(Constants.REPOSITORIES_MAP + Configurations.SEPARATOR + group.getName()); http://git-wip-us.apache.org/repos/asf/karaf-cellar/blob/7c3e1c70/features/src/main/resources/OSGI-INF/blueprint/blueprint.xml ---------------------------------------------------------------------- diff --git a/features/src/main/resources/OSGI-INF/blueprint/blueprint.xml b/features/src/main/resources/OSGI-INF/blueprint/blueprint.xml index 4f30536..7e319ea 100644 --- a/features/src/main/resources/OSGI-INF/blueprint/blueprint.xml +++ b/features/src/main/resources/OSGI-INF/blueprint/blueprint.xml @@ -34,6 +34,7 @@ <property name="clusterManager" ref="clusterManager"/> <property name="groupManager" ref="groupManager"/> <property name="configurationAdmin" ref="configurationAdmin"/> + <property name="eventProducer" ref="eventProducer"/> <property name="featuresService" ref="featuresService"/> </bean> <service ref="synchronizer" interface="org.apache.karaf.cellar.core.Synchronizer">
