Author: cziegeler
Date: Thu Apr 30 13:17:56 2015
New Revision: 1676979

URL: http://svn.apache.org/r1676979
Log:
SLING-4683 : Ensure that topology change events are processed quickly

Added:
    
sling/trunk/bundles/extensions/event/src/main/java/org/apache/sling/event/impl/jobs/config/TopologyHandler.java
   (with props)
Modified:
    
sling/trunk/bundles/extensions/event/src/main/java/org/apache/sling/event/impl/jobs/config/JobManagerConfiguration.java
    
sling/trunk/bundles/extensions/event/src/main/java/org/apache/sling/event/impl/jobs/timed/TimedEventSender.java
    
sling/trunk/bundles/extensions/event/src/test/java/org/apache/sling/event/it/ChaosTest.java

Modified: 
sling/trunk/bundles/extensions/event/src/main/java/org/apache/sling/event/impl/jobs/config/JobManagerConfiguration.java
URL: 
http://svn.apache.org/viewvc/sling/trunk/bundles/extensions/event/src/main/java/org/apache/sling/event/impl/jobs/config/JobManagerConfiguration.java?rev=1676979&r1=1676978&r2=1676979&view=diff
==============================================================================
--- 
sling/trunk/bundles/extensions/event/src/main/java/org/apache/sling/event/impl/jobs/config/JobManagerConfiguration.java
 (original)
+++ 
sling/trunk/bundles/extensions/event/src/main/java/org/apache/sling/event/impl/jobs/config/JobManagerConfiguration.java
 Thu Apr 30 13:17:56 2015
@@ -63,7 +63,7 @@ import org.slf4j.LoggerFactory;
            label="Apache Sling Job Manager",
            description="This is the central service of the job handling.",
            name="org.apache.sling.event.impl.jobs.jcr.PersistenceHandler")
-@Service(value={JobManagerConfiguration.class, TopologyEventListener.class})
+@Service(value={JobManagerConfiguration.class})
 @Properties({
     @Property(name=JobManagerConfiguration.PROPERTY_DISABLE_DISTRIBUTION,
               boolValue=JobManagerConfiguration.DEFAULT_DISABLE_DISTRIBUTION,

Added: 
sling/trunk/bundles/extensions/event/src/main/java/org/apache/sling/event/impl/jobs/config/TopologyHandler.java
URL: 
http://svn.apache.org/viewvc/sling/trunk/bundles/extensions/event/src/main/java/org/apache/sling/event/impl/jobs/config/TopologyHandler.java?rev=1676979&view=auto
==============================================================================
--- 
sling/trunk/bundles/extensions/event/src/main/java/org/apache/sling/event/impl/jobs/config/TopologyHandler.java
 (added)
+++ 
sling/trunk/bundles/extensions/event/src/main/java/org/apache/sling/event/impl/jobs/config/TopologyHandler.java
 Thu Apr 30 13:17:56 2015
@@ -0,0 +1,109 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.sling.event.impl.jobs.config;
+
+import java.util.concurrent.BlockingQueue;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+import org.apache.felix.scr.annotations.Activate;
+import org.apache.felix.scr.annotations.Component;
+import org.apache.felix.scr.annotations.Deactivate;
+import org.apache.felix.scr.annotations.Reference;
+import org.apache.felix.scr.annotations.Service;
+import org.apache.sling.discovery.TopologyEvent;
+import org.apache.sling.discovery.TopologyEventListener;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+@Component
+@Service(value = TopologyEventListener.class)
+public class TopologyHandler implements TopologyEventListener, Runnable {
+
+    /** The logger. */
+    private final Logger logger = 
LoggerFactory.getLogger(this.getClass().getName());
+
+    @Reference
+    private JobManagerConfiguration configuration;
+
+    /** A local queue for async handling of the events */
+    private final BlockingQueue<QueueItem> queue = new 
LinkedBlockingQueue<QueueItem>();
+
+    /** Active flag. */
+    private final AtomicBoolean isActive = new AtomicBoolean(false);
+
+    @Activate
+    protected void activate() {
+        this.isActive.set(true);
+        final Thread thread = new Thread(this, "Apache Sling Job Topology 
Listener Thread");
+        thread.setDaemon(true);
+
+        thread.start();
+    }
+
+    @Deactivate
+    protected void deactivate() {
+        this.isActive.set(false);
+        try {
+            this.queue.put(new QueueItem());
+        } catch ( final InterruptedException ie) {
+            logger.warn("Thread got interrupted.", ie);
+            Thread.currentThread().interrupt();
+        }
+    }
+
+    @Override
+    public void handleTopologyEvent(final TopologyEvent event) {
+        final QueueItem item = new QueueItem();
+        item.event = event;
+        try {
+            this.queue.put(item);
+        } catch ( final InterruptedException ie) {
+            logger.warn("Thread got interrupted.", ie);
+            Thread.currentThread().interrupt();
+        }
+    }
+
+    @Override
+    public void run() {
+        while ( isActive.get() ) {
+            QueueItem item = null;
+            try {
+                item = this.queue.take();
+            } catch ( final InterruptedException ie) {
+                logger.warn("Thread got interrupted.", ie);
+                Thread.currentThread().interrupt();
+                isActive.set(false);
+            }
+            if ( isActive.get() && item != null && item.event != null ) {
+                final JobManagerConfiguration config = this.configuration;
+                if ( config != null ) {
+                    config.handleTopologyEvent(item.event);
+                }
+            }
+        }
+    }
+
+    /**
+     * We need a holder class to be able to put something into the queue to 
stop it.
+     */
+    public static final class QueueItem {
+        public TopologyEvent event;
+    }
+}

Propchange: 
sling/trunk/bundles/extensions/event/src/main/java/org/apache/sling/event/impl/jobs/config/TopologyHandler.java
------------------------------------------------------------------------------
    svn:eol-style = native

Propchange: 
sling/trunk/bundles/extensions/event/src/main/java/org/apache/sling/event/impl/jobs/config/TopologyHandler.java
------------------------------------------------------------------------------
    svn:keywords = author date id revision rev url

Propchange: 
sling/trunk/bundles/extensions/event/src/main/java/org/apache/sling/event/impl/jobs/config/TopologyHandler.java
------------------------------------------------------------------------------
    svn:mime-type = text/plain

Modified: 
sling/trunk/bundles/extensions/event/src/main/java/org/apache/sling/event/impl/jobs/timed/TimedEventSender.java
URL: 
http://svn.apache.org/viewvc/sling/trunk/bundles/extensions/event/src/main/java/org/apache/sling/event/impl/jobs/timed/TimedEventSender.java?rev=1676979&r1=1676978&r2=1676979&view=diff
==============================================================================
--- 
sling/trunk/bundles/extensions/event/src/main/java/org/apache/sling/event/impl/jobs/timed/TimedEventSender.java
 (original)
+++ 
sling/trunk/bundles/extensions/event/src/main/java/org/apache/sling/event/impl/jobs/timed/TimedEventSender.java
 Thu Apr 30 13:17:56 2015
@@ -86,7 +86,7 @@ import org.slf4j.LoggerFactory;
                  ResourceHelper.BUNDLE_EVENT_STARTED,
                  ResourceHelper.BUNDLE_EVENT_UPDATED})
 public class TimedEventSender
-    implements Job, TimedEventStatusProvider, EventHandler, 
TopologyEventListener {
+    implements Job, TimedEventStatusProvider, EventHandler, 
TopologyEventListener, Runnable {
 
     private static final String JOB_TOPIC = "topic";
 
@@ -126,12 +126,19 @@ public class TimedEventSender
 
     private final AtomicBoolean threadStarted = new AtomicBoolean(false);
 
+    /** A local queue for async handling of the topology events */
+    private final BlockingQueue<QueueItem> topologyEventQueue = new 
LinkedBlockingQueue<QueueItem>();
+
     /**
      * Activate this component.
      */
     @Activate
     protected void activate() {
         this.running = true;
+        final Thread thread = new Thread(this, "Apache Sling Timed Job 
Topology Listener Thread");
+        thread.setDaemon(true);
+
+        thread.start();
     }
 
     /**
@@ -141,6 +148,12 @@ public class TimedEventSender
     protected void deactivate() {
         this.running = false;
         this.stopScheduling();
+        try {
+            this.topologyEventQueue.put(new QueueItem());
+        } catch ( final InterruptedException ie) {
+            logger.warn("Thread got interrupted.", ie);
+            Thread.currentThread().interrupt();
+        }
     }
 
     private void stopScheduling() {
@@ -718,23 +731,51 @@ public class TimedEventSender
         }
     }
 
-    /**
-     * @see 
org.apache.sling.discovery.TopologyEventListener#handleTopologyEvent(org.apache.sling.discovery.TopologyEvent)
-     */
     @Override
     public void handleTopologyEvent(final TopologyEvent event) {
-        if ( event.getType() == Type.TOPOLOGY_CHANGING ) {
-            this.active = false;
-            this.stopScheduling();
-        } else if ( event.getType() == Type.TOPOLOGY_CHANGED || 
event.getType() == Type.TOPOLOGY_INIT ) {
-            final boolean previouslyActive = this.active;
-            this.active = event.getNewView().getLocalInstance().isLeader();
-            if ( this.active && !previouslyActive ) {
-                this.startScheduling();
+        final QueueItem item = new QueueItem();
+        item.event = event;
+        try {
+            this.topologyEventQueue.put(item);
+        } catch ( final InterruptedException ie) {
+            logger.warn("Thread got interrupted.", ie);
+            Thread.currentThread().interrupt();
+        }
+    }
+
+    @Override
+    public void run() {
+        while ( this.running ) {
+            QueueItem item = null;
+            try {
+                item = this.topologyEventQueue.take();
+            } catch ( final InterruptedException ie) {
+                logger.warn("Thread got interrupted.", ie);
+                Thread.currentThread().interrupt();
+                this.running = false;
             }
-            if ( !this.active && previouslyActive ) {
-                this.stopScheduling();
+            if ( this.running && item != null && item.event != null ) {
+                if ( item.event.getType() == Type.TOPOLOGY_CHANGING ) {
+                    this.active = false;
+                    this.stopScheduling();
+                } else if ( item.event.getType() == Type.TOPOLOGY_CHANGED || 
item.event.getType() == Type.TOPOLOGY_INIT ) {
+                    final boolean previouslyActive = this.active;
+                    this.active = 
item.event.getNewView().getLocalInstance().isLeader();
+                    if ( this.active && !previouslyActive ) {
+                        this.startScheduling();
+                    }
+                    if ( !this.active && previouslyActive ) {
+                        this.stopScheduling();
+                    }
+                }
             }
         }
     }
+
+    /**
+     * We need a holder class to be able to put something into the queue to 
stop it.
+     */
+    public static final class QueueItem {
+        public TopologyEvent event;
+    }
 }

Modified: 
sling/trunk/bundles/extensions/event/src/test/java/org/apache/sling/event/it/ChaosTest.java
URL: 
http://svn.apache.org/viewvc/sling/trunk/bundles/extensions/event/src/test/java/org/apache/sling/event/it/ChaosTest.java?rev=1676979&r1=1676978&r2=1676979&view=diff
==============================================================================
--- 
sling/trunk/bundles/extensions/event/src/test/java/org/apache/sling/event/it/ChaosTest.java
 (original)
+++ 
sling/trunk/bundles/extensions/event/src/test/java/org/apache/sling/event/it/ChaosTest.java
 Thu Apr 30 13:17:56 2015
@@ -18,8 +18,8 @@
  */
 package org.apache.sling.event.it;
 
-import static org.junit.Assert.assertEquals;
 import static org.junit.Assert.assertNotNull;
+import static org.junit.Assert.assertTrue;
 
 import java.io.IOException;
 import java.util.ArrayList;
@@ -268,11 +268,22 @@ public class ChaosTest extends AbstractJ
         final TopologyView view = views.get(0);
 
         try {
-            final ServiceReference[] refs = 
this.bc.getServiceReferences(TopologyEventListener.class.getName(),
-                    
"(objectClass=org.apache.sling.event.impl.jobs.config.JobManagerConfiguration)");
+            final ServiceReference[] refs = 
this.bc.getServiceReferences(TopologyEventListener.class.getName(), null);
             assertNotNull(refs);
-            assertEquals(1, refs.length);
-            final TopologyEventListener tel = 
(TopologyEventListener)bc.getService(refs[0]);
+            assertTrue(refs.length > 1);
+            int index = 0;
+            TopologyEventListener found = null;
+            while ( index < refs.length ) {
+                final TopologyEventListener listener = (TopologyEventListener) 
this.bc.getService(refs[index]);
+                if ( 
listener.getClass().getName().equals("org.apache.sling.event.impl.jobs.config.TopologyHandler")
 ) {
+                    found = listener;
+                    break;
+                }
+                bc.ungetService(refs[index]);
+                index++;
+            }
+            assertNotNull(found);
+            final TopologyEventListener tel = found;
 
             threads.add(new Thread() {
 


Reply via email to