Author: catholicon
Date: Fri Dec  8 14:02:32 2017
New Revision: 1817497

URL: http://svn.apache.org/viewvc?rev=1817497&view=rev
Log:
OAK-7042: Pin async indexer on cluster leader

Add a method to WhiteboardUtils to allow scheduling on leader. Use that
to setup async indexer to be pinned to leader

Modified:
    
jackrabbit/oak/trunk/oak-core-spi/src/main/java/org/apache/jackrabbit/oak/spi/whiteboard/WhiteboardUtils.java
    
jackrabbit/oak/trunk/oak-core/src/main/java/org/apache/jackrabbit/oak/plugins/index/IndexMBeanRegistration.java
    
jackrabbit/oak/trunk/oak-core/src/test/java/org/apache/jackrabbit/oak/plugins/index/IndexMBeanRegistrationTest.java
    
jackrabbit/oak/trunk/oak-core/src/test/java/org/apache/jackrabbit/oak/spi/whiteboard/WhiteboardUtilsTest.java

Modified: 
jackrabbit/oak/trunk/oak-core-spi/src/main/java/org/apache/jackrabbit/oak/spi/whiteboard/WhiteboardUtils.java
URL: 
http://svn.apache.org/viewvc/jackrabbit/oak/trunk/oak-core-spi/src/main/java/org/apache/jackrabbit/oak/spi/whiteboard/WhiteboardUtils.java?rev=1817497&r1=1817496&r2=1817497&view=diff
==============================================================================
--- 
jackrabbit/oak/trunk/oak-core-spi/src/main/java/org/apache/jackrabbit/oak/spi/whiteboard/WhiteboardUtils.java
 (original)
+++ 
jackrabbit/oak/trunk/oak-core-spi/src/main/java/org/apache/jackrabbit/oak/spi/whiteboard/WhiteboardUtils.java
 Fri Dec  8 14:02:32 2017
@@ -34,6 +34,10 @@ import com.google.common.collect.Immutab
 import com.google.common.collect.ImmutableMap;
 import com.google.common.collect.Iterables;
 
+import static 
org.apache.jackrabbit.oak.spi.whiteboard.WhiteboardUtils.ScheduleExecutionInstanceTypes.DEFAULT;
+import static 
org.apache.jackrabbit.oak.spi.whiteboard.WhiteboardUtils.ScheduleExecutionInstanceTypes.RUN_ON_LEADER;
+import static 
org.apache.jackrabbit.oak.spi.whiteboard.WhiteboardUtils.ScheduleExecutionInstanceTypes.RUN_ON_SINGLE;
+
 public class WhiteboardUtils {
 
     /**
@@ -41,6 +45,12 @@ public class WhiteboardUtils {
      */
     public static final String JMX_OAK_DOMAIN = "org.apache.jackrabbit.oak";
 
+    public enum ScheduleExecutionInstanceTypes {
+        DEFAULT,
+        RUN_ON_SINGLE,
+        RUN_ON_LEADER
+    }
+
     public static Registration scheduleWithFixedDelay(
             Whiteboard whiteboard, Runnable runnable, long delayInSeconds) {
         return scheduleWithFixedDelay(whiteboard, runnable, delayInSeconds, 
false, false);
@@ -56,13 +66,25 @@ public class WhiteboardUtils {
     public static Registration scheduleWithFixedDelay(
             Whiteboard whiteboard, Runnable runnable, Map<String, Object> 
extraProps, long delayInSeconds, boolean runOnSingleClusterNode,
             boolean useDedicatedPool) {
+        return scheduleWithFixedDelay(whiteboard, runnable, extraProps, 
delayInSeconds,
+                runOnSingleClusterNode ? RUN_ON_SINGLE : DEFAULT,
+                useDedicatedPool);
+    }
+
+    public static Registration scheduleWithFixedDelay(
+            Whiteboard whiteboard, Runnable runnable, Map<String, Object> 
extraProps, long delayInSeconds,
+            ScheduleExecutionInstanceTypes scheduleExecutionInstanceTypes, 
boolean useDedicatedPool) {
+
         ImmutableMap.Builder<String, Object> builder = ImmutableMap.<String, 
Object>builder()
                 .putAll(extraProps)
                 .put("scheduler.period", delayInSeconds)
                 .put("scheduler.concurrent", false);
-        if (runOnSingleClusterNode) {
-            //Make use of feature while running in Sling SLING-2979
+        if (scheduleExecutionInstanceTypes == RUN_ON_SINGLE) {
+            //Make use of feature while running in Sling SLING-5387
             builder.put("scheduler.runOn", "SINGLE");
+        } else if (scheduleExecutionInstanceTypes == RUN_ON_LEADER) {
+            //Make use of feature while running in Sling SLING-2979
+            builder.put("scheduler.runOn", "LEADER");
         }
         if (useDedicatedPool) {
             //Make use of dedicated threadpool SLING-5831

Modified: 
jackrabbit/oak/trunk/oak-core/src/main/java/org/apache/jackrabbit/oak/plugins/index/IndexMBeanRegistration.java
URL: 
http://svn.apache.org/viewvc/jackrabbit/oak/trunk/oak-core/src/main/java/org/apache/jackrabbit/oak/plugins/index/IndexMBeanRegistration.java?rev=1817497&r1=1817496&r2=1817497&view=diff
==============================================================================
--- 
jackrabbit/oak/trunk/oak-core/src/main/java/org/apache/jackrabbit/oak/plugins/index/IndexMBeanRegistration.java
 (original)
+++ 
jackrabbit/oak/trunk/oak-core/src/main/java/org/apache/jackrabbit/oak/plugins/index/IndexMBeanRegistration.java
 Fri Dec  8 14:02:32 2017
@@ -16,6 +16,7 @@
  */
 package org.apache.jackrabbit.oak.plugins.index;
 
+import static 
org.apache.jackrabbit.oak.spi.whiteboard.WhiteboardUtils.ScheduleExecutionInstanceTypes.RUN_ON_LEADER;
 import static 
org.apache.jackrabbit.oak.spi.whiteboard.WhiteboardUtils.registerMBean;
 import static 
org.apache.jackrabbit.oak.spi.whiteboard.WhiteboardUtils.scheduleWithFixedDelay;
 
@@ -29,6 +30,7 @@ import org.apache.jackrabbit.oak.spi.whi
 import org.apache.jackrabbit.oak.spi.whiteboard.Whiteboard;
 
 import com.google.common.collect.Lists;
+import org.apache.jackrabbit.oak.spi.whiteboard.WhiteboardUtils;
 
 public class IndexMBeanRegistration implements Registration {
 
@@ -45,7 +47,7 @@ public class IndexMBeanRegistration impl
                 AsyncIndexUpdate.PROP_ASYNC_NAME, task.getName(),
                 "scheduler.name", AsyncIndexUpdate.class.getName() + "-" + 
task.getName()
         );
-        regs.add(scheduleWithFixedDelay(whiteboard, task, config, 
delayInSeconds, true, true));
+        regs.add(scheduleWithFixedDelay(whiteboard, task, config, 
delayInSeconds, RUN_ON_LEADER, true));
         regs.add(registerMBean(whiteboard, IndexStatsMBean.class,
                 task.getIndexStats(), IndexStatsMBean.TYPE, task.getName()));
     }

Modified: 
jackrabbit/oak/trunk/oak-core/src/test/java/org/apache/jackrabbit/oak/plugins/index/IndexMBeanRegistrationTest.java
URL: 
http://svn.apache.org/viewvc/jackrabbit/oak/trunk/oak-core/src/test/java/org/apache/jackrabbit/oak/plugins/index/IndexMBeanRegistrationTest.java?rev=1817497&r1=1817496&r2=1817497&view=diff
==============================================================================
--- 
jackrabbit/oak/trunk/oak-core/src/test/java/org/apache/jackrabbit/oak/plugins/index/IndexMBeanRegistrationTest.java
 (original)
+++ 
jackrabbit/oak/trunk/oak-core/src/test/java/org/apache/jackrabbit/oak/plugins/index/IndexMBeanRegistrationTest.java
 Fri Dec  8 14:02:32 2017
@@ -42,16 +42,25 @@ public class IndexMBeanRegistrationTest
                 return super.register(type, service, properties);
             }
         };
+        long schedulingDelayInSecs = 7; // some number which is hard to 
default else-where
 
         AsyncIndexUpdate update = new AsyncIndexUpdate("async",
                 new MemoryNodeStore(), new CompositeIndexEditorProvider());
         IndexMBeanRegistration reg = new IndexMBeanRegistration(wb);
-        reg.registerAsyncIndexer(update, 5);
+        reg.registerAsyncIndexer(update, schedulingDelayInSecs);
         try {
             Map<?, ?> map = props.get();
             assertNotNull(map);
             assertEquals(AsyncIndexUpdate.class.getName() + "-async",
                     map.get("scheduler.name"));
+            assertEquals(schedulingDelayInSecs,
+                    map.get("scheduler.period"));
+            assertEquals(false,
+                    map.get("scheduler.concurrent"));
+            assertEquals("LEADER",
+                    map.get("scheduler.runOn"));
+            assertEquals("oak",
+                    map.get("scheduler.threadPool"));
         } finally {
             reg.unregister();
         }

Modified: 
jackrabbit/oak/trunk/oak-core/src/test/java/org/apache/jackrabbit/oak/spi/whiteboard/WhiteboardUtilsTest.java
URL: 
http://svn.apache.org/viewvc/jackrabbit/oak/trunk/oak-core/src/test/java/org/apache/jackrabbit/oak/spi/whiteboard/WhiteboardUtilsTest.java?rev=1817497&r1=1817496&r2=1817497&view=diff
==============================================================================
--- 
jackrabbit/oak/trunk/oak-core/src/test/java/org/apache/jackrabbit/oak/spi/whiteboard/WhiteboardUtilsTest.java
 (original)
+++ 
jackrabbit/oak/trunk/oak-core/src/test/java/org/apache/jackrabbit/oak/spi/whiteboard/WhiteboardUtilsTest.java
 Fri Dec  8 14:02:32 2017
@@ -20,6 +20,7 @@
 package org.apache.jackrabbit.oak.spi.whiteboard;
 
 import java.lang.management.ManagementFactory;
+import java.util.Collections;
 import java.util.List;
 import java.util.Map;
 import java.util.concurrent.atomic.AtomicReference;
@@ -35,6 +36,9 @@ import org.apache.jackrabbit.oak.query.Q
 import org.junit.After;
 import org.junit.Test;
 
+import static 
org.apache.jackrabbit.oak.spi.whiteboard.WhiteboardUtils.ScheduleExecutionInstanceTypes.DEFAULT;
+import static 
org.apache.jackrabbit.oak.spi.whiteboard.WhiteboardUtils.ScheduleExecutionInstanceTypes.RUN_ON_LEADER;
+import static 
org.apache.jackrabbit.oak.spi.whiteboard.WhiteboardUtils.ScheduleExecutionInstanceTypes.RUN_ON_SINGLE;
 import static org.junit.Assert.assertEquals;
 import static org.junit.Assert.assertNotNull;
 import static org.junit.Assert.assertNull;
@@ -127,6 +131,69 @@ public class WhiteboardUtilsTest {
         assertEquals("bar", props.get().get("foo"));
     }
 
+    @Test
+    public void scheduledJobDefaultExecutionInstanceType() {
+        Map<String, Object> config = Collections.emptyMap();
+        final AtomicReference<Map<?, ?>> props = new AtomicReference<Map<?, 
?>>();
+        Whiteboard wb = new DefaultWhiteboard(){
+            @Override
+            public <T> Registration register(Class<T> type, T service, Map<?, 
?> properties) {
+                props.set(properties);
+                return super.register(type, service, properties);
+            }
+        };
+
+        WhiteboardUtils.scheduleWithFixedDelay(wb, new TestRunnable(), 1);
+        assertNull(props.get().get("scheduler.runOn"));
+
+        WhiteboardUtils.scheduleWithFixedDelay(wb, new TestRunnable(), 1, 
false, false);
+        assertNull(props.get().get("scheduler.runOn"));
+
+        WhiteboardUtils.scheduleWithFixedDelay(wb, new TestRunnable(), 
config,1, false, false);
+        assertNull(props.get().get("scheduler.runOn"));
+
+        WhiteboardUtils.scheduleWithFixedDelay(wb, new TestRunnable(), config, 
1, DEFAULT, false);
+        assertNull(props.get().get("scheduler.runOn"));
+    }
+
+    @Test
+    public void scheduledJobOnSingle() {
+        Map<String, Object> config = Collections.emptyMap();
+        final AtomicReference<Map<?, ?>> props = new AtomicReference<Map<?, 
?>>();
+        Whiteboard wb = new DefaultWhiteboard(){
+            @Override
+            public <T> Registration register(Class<T> type, T service, Map<?, 
?> properties) {
+                props.set(properties);
+                return super.register(type, service, properties);
+            }
+        };
+
+        WhiteboardUtils.scheduleWithFixedDelay(wb, new TestRunnable(), 1, 
true, false);
+        assertEquals("SINGLE", props.get().get("scheduler.runOn"));
+
+        WhiteboardUtils.scheduleWithFixedDelay(wb, new TestRunnable(), config, 
1, true, false);
+        assertEquals("SINGLE", props.get().get("scheduler.runOn"));
+
+        WhiteboardUtils.scheduleWithFixedDelay(wb, new TestRunnable(), config, 
1, RUN_ON_SINGLE, false);
+        assertEquals("SINGLE", props.get().get("scheduler.runOn"));
+    }
+
+    @Test
+    public void scheduledJobOnLeader() {
+        Map<String, Object> config = Collections.emptyMap();
+        final AtomicReference<Map<?, ?>> props = new AtomicReference<Map<?, 
?>>();
+        Whiteboard wb = new DefaultWhiteboard(){
+            @Override
+            public <T> Registration register(Class<T> type, T service, Map<?, 
?> properties) {
+                props.set(properties);
+                return super.register(type, service, properties);
+            }
+        };
+
+        WhiteboardUtils.scheduleWithFixedDelay(wb, new TestRunnable(), config, 
1, RUN_ON_LEADER, false);
+        assertEquals("LEADER", props.get().get("scheduler.runOn"));
+    }
+
     public interface HelloMBean {
         boolean isRunning();
         int getCount();


Reply via email to