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">

Reply via email to