You are viewing a plain text version of this content. The canonical link for it is here.
Posted to commits@camel.apache.org by bv...@apache.org on 2012/02/26 16:30:00 UTC

svn commit: r1293852 - in /camel/trunk/camel-core/src: main/java/org/apache/camel/api/management/mbean/ main/java/org/apache/camel/management/ main/java/org/apache/camel/management/mbean/ main/java/org/apache/camel/processor/idempotent/ test/java/org/a...

Author: bvahdat
Date: Sun Feb 26 15:29:59 2012
New Revision: 1293852

URL: http://svn.apache.org/viewvc?rev=1293852&view=rev
Log:
CAMEL-4782: Add ManagedIdempotentConsumerBean providing get & reset of the duplicate messages count detected.

Added:
    camel/trunk/camel-core/src/main/java/org/apache/camel/api/management/mbean/ManagedIdempotentConsumerMBean.java   (with props)
    camel/trunk/camel-core/src/main/java/org/apache/camel/management/mbean/ManagedIdempotentConsumer.java   (with props)
Modified:
    camel/trunk/camel-core/src/main/java/org/apache/camel/management/DefaultManagementLifecycleStrategy.java
    camel/trunk/camel-core/src/main/java/org/apache/camel/management/DefaultManagementObjectStrategy.java
    camel/trunk/camel-core/src/main/java/org/apache/camel/processor/idempotent/IdempotentConsumer.java
    camel/trunk/camel-core/src/test/java/org/apache/camel/management/ManagedMemoryIdempotentConsumerTest.java

Added: camel/trunk/camel-core/src/main/java/org/apache/camel/api/management/mbean/ManagedIdempotentConsumerMBean.java
URL: http://svn.apache.org/viewvc/camel/trunk/camel-core/src/main/java/org/apache/camel/api/management/mbean/ManagedIdempotentConsumerMBean.java?rev=1293852&view=auto
==============================================================================
--- camel/trunk/camel-core/src/main/java/org/apache/camel/api/management/mbean/ManagedIdempotentConsumerMBean.java (added)
+++ camel/trunk/camel-core/src/main/java/org/apache/camel/api/management/mbean/ManagedIdempotentConsumerMBean.java Sun Feb 26 15:29:59 2012
@@ -0,0 +1,30 @@
+/**
+ * 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.camel.api.management.mbean;
+
+import org.apache.camel.api.management.ManagedAttribute;
+import org.apache.camel.api.management.ManagedOperation;
+
+public interface ManagedIdempotentConsumerMBean extends ManagedProcessorMBean {
+
+    @ManagedAttribute(description = "Current count of duplicate Messages")
+    long getDuplicateMessageCount();
+
+    @ManagedOperation(description = "Reset the current count of duplicate Messages")
+    void resetDuplicateMessageCount();
+
+}

Propchange: camel/trunk/camel-core/src/main/java/org/apache/camel/api/management/mbean/ManagedIdempotentConsumerMBean.java
------------------------------------------------------------------------------
    svn:eol-style = native

Modified: camel/trunk/camel-core/src/main/java/org/apache/camel/management/DefaultManagementLifecycleStrategy.java
URL: http://svn.apache.org/viewvc/camel/trunk/camel-core/src/main/java/org/apache/camel/management/DefaultManagementLifecycleStrategy.java?rev=1293852&r1=1293851&r2=1293852&view=diff
==============================================================================
--- camel/trunk/camel-core/src/main/java/org/apache/camel/management/DefaultManagementLifecycleStrategy.java (original)
+++ camel/trunk/camel-core/src/main/java/org/apache/camel/management/DefaultManagementLifecycleStrategy.java Sun Feb 26 15:29:59 2012
@@ -422,10 +422,9 @@ public class DefaultManagementLifecycleS
             ManagedService ms = (ManagedService) answer;
             ms.setRoute(route);
             ms.init(getManagementStrategy());
-            return answer;
-        } else {
-            return answer;
         }
+
+        return answer;
     }
 
     private Object getManagedObjectForProcessor(CamelContext context, Processor processor, Route route) {

Modified: camel/trunk/camel-core/src/main/java/org/apache/camel/management/DefaultManagementObjectStrategy.java
URL: http://svn.apache.org/viewvc/camel/trunk/camel-core/src/main/java/org/apache/camel/management/DefaultManagementObjectStrategy.java?rev=1293852&r1=1293851&r2=1293852&view=diff
==============================================================================
--- camel/trunk/camel-core/src/main/java/org/apache/camel/management/DefaultManagementObjectStrategy.java (original)
+++ camel/trunk/camel-core/src/main/java/org/apache/camel/management/DefaultManagementObjectStrategy.java Sun Feb 26 15:29:59 2012
@@ -39,6 +39,7 @@ import org.apache.camel.management.mbean
 import org.apache.camel.management.mbean.ManagedEndpoint;
 import org.apache.camel.management.mbean.ManagedErrorHandler;
 import org.apache.camel.management.mbean.ManagedEventNotifier;
+import org.apache.camel.management.mbean.ManagedIdempotentConsumer;
 import org.apache.camel.management.mbean.ManagedProcessor;
 import org.apache.camel.management.mbean.ManagedProducer;
 import org.apache.camel.management.mbean.ManagedRoute;
@@ -54,6 +55,7 @@ import org.apache.camel.processor.Delaye
 import org.apache.camel.processor.ErrorHandler;
 import org.apache.camel.processor.SendProcessor;
 import org.apache.camel.processor.Throttler;
+import org.apache.camel.processor.idempotent.IdempotentConsumer;
 import org.apache.camel.spi.BrowsableEndpoint;
 import org.apache.camel.spi.EventNotifier;
 import org.apache.camel.spi.ManagementObjectStrategy;
@@ -178,6 +180,8 @@ public class DefaultManagementObjectStra
                 answer = new ManagedSendProcessor(context, (SendProcessor) target, definition);
             } else if (target instanceof BeanProcessor) {
                 answer = new ManagedBeanProcessor(context, (BeanProcessor) target, definition);
+            } else if (target instanceof IdempotentConsumer) {
+                answer = new ManagedIdempotentConsumer(context, (IdempotentConsumer) target, definition);
             } else if (target instanceof org.apache.camel.spi.ManagementAware) {
                 return ((org.apache.camel.spi.ManagementAware<Processor>) target).getManagedObject(processor);
             }

Added: camel/trunk/camel-core/src/main/java/org/apache/camel/management/mbean/ManagedIdempotentConsumer.java
URL: http://svn.apache.org/viewvc/camel/trunk/camel-core/src/main/java/org/apache/camel/management/mbean/ManagedIdempotentConsumer.java?rev=1293852&view=auto
==============================================================================
--- camel/trunk/camel-core/src/main/java/org/apache/camel/management/mbean/ManagedIdempotentConsumer.java (added)
+++ camel/trunk/camel-core/src/main/java/org/apache/camel/management/mbean/ManagedIdempotentConsumer.java Sun Feb 26 15:29:59 2012
@@ -0,0 +1,47 @@
+/**
+ * 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.camel.management.mbean;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.api.management.ManagedResource;
+import org.apache.camel.api.management.mbean.ManagedIdempotentConsumerMBean;
+import org.apache.camel.model.ProcessorDefinition;
+import org.apache.camel.processor.idempotent.IdempotentConsumer;
+
+@ManagedResource(description = "Managed Idempotent Consumer")
+public class ManagedIdempotentConsumer extends ManagedProcessor implements ManagedIdempotentConsumerMBean {
+
+    public ManagedIdempotentConsumer(CamelContext context, IdempotentConsumer idempotentConsumer, ProcessorDefinition<?> definition) {
+        super(context, idempotentConsumer, definition);
+    }
+
+    @Override
+    public IdempotentConsumer getProcessor() {
+        return (IdempotentConsumer) super.getProcessor();
+    }
+
+    @Override
+    public long getDuplicateMessageCount() {
+        return getProcessor().getDuplicateMessageCount();
+    }
+
+    @Override
+    public void resetDuplicateMessageCount() {
+        getProcessor().resetDuplicateMessageCount();
+    }
+
+}

Propchange: camel/trunk/camel-core/src/main/java/org/apache/camel/management/mbean/ManagedIdempotentConsumer.java
------------------------------------------------------------------------------
    svn:eol-style = native

Modified: camel/trunk/camel-core/src/main/java/org/apache/camel/processor/idempotent/IdempotentConsumer.java
URL: http://svn.apache.org/viewvc/camel/trunk/camel-core/src/main/java/org/apache/camel/processor/idempotent/IdempotentConsumer.java?rev=1293852&r1=1293851&r2=1293852&view=diff
==============================================================================
--- camel/trunk/camel-core/src/main/java/org/apache/camel/processor/idempotent/IdempotentConsumer.java (original)
+++ camel/trunk/camel-core/src/main/java/org/apache/camel/processor/idempotent/IdempotentConsumer.java Sun Feb 26 15:29:59 2012
@@ -18,6 +18,7 @@ package org.apache.camel.processor.idemp
 
 import java.util.ArrayList;
 import java.util.List;
+import java.util.concurrent.atomic.AtomicLong;
 
 import org.apache.camel.AsyncCallback;
 import org.apache.camel.AsyncProcessor;
@@ -45,6 +46,7 @@ public class IdempotentConsumer extends 
     private final boolean eager;
     private final boolean skipDuplicate;
     private final boolean removeOnFailure;
+    private final AtomicLong duplicateMessageCount = new AtomicLong();
 
     public IdempotentConsumer(Expression messageIdExpression, IdempotentRepository<String> idempotentRepository,
                               boolean eager, boolean skipDuplicate, boolean removeOnFailure, Processor processor) {
@@ -84,11 +86,10 @@ public class IdempotentConsumer extends 
         if (!newKey) {
             // mark the exchange as duplicate
             exchange.setProperty(Exchange.DUPLICATE_MESSAGE, Boolean.TRUE);
-        }
 
-        if (!newKey) {
             // we already have this key so its a duplicate message
-            onDuplicateMessage(exchange, messageId);
+            onDuplicate(exchange, messageId);
+
             if (skipDuplicate) {
                 // if we should skip duplicate then we are done
                 LOG.debug("Ignoring duplicate message with id: {} for exchange: {}", messageId, exchange);
@@ -131,6 +132,10 @@ public class IdempotentConsumer extends 
         return processor;
     }
 
+    public long getDuplicateMessageCount() {
+        return duplicateMessageCount.get();
+    }
+
     // Implementation methods
     // -------------------------------------------------------------------------
 
@@ -143,6 +148,19 @@ public class IdempotentConsumer extends 
     }
 
     /**
+     * Resets the duplicate message counter to <code>0L</code>.
+     */
+    public void resetDuplicateMessageCount() {
+        duplicateMessageCount.set(0L);
+    }
+
+    private void onDuplicate(Exchange exchange, String messageId) {
+        duplicateMessageCount.incrementAndGet();
+
+        onDuplicateMessage(exchange, messageId);
+    }
+
+    /**
      * A strategy method to allow derived classes to overload the behaviour of
      * processing a duplicate message
      *

Modified: camel/trunk/camel-core/src/test/java/org/apache/camel/management/ManagedMemoryIdempotentConsumerTest.java
URL: http://svn.apache.org/viewvc/camel/trunk/camel-core/src/test/java/org/apache/camel/management/ManagedMemoryIdempotentConsumerTest.java?rev=1293852&r1=1293851&r2=1293852&view=diff
==============================================================================
--- camel/trunk/camel-core/src/test/java/org/apache/camel/management/ManagedMemoryIdempotentConsumerTest.java (original)
+++ camel/trunk/camel-core/src/test/java/org/apache/camel/management/ManagedMemoryIdempotentConsumerTest.java Sun Feb 26 15:29:59 2012
@@ -28,7 +28,6 @@ import org.apache.camel.builder.RouteBui
 import org.apache.camel.component.mock.MockEndpoint;
 import org.apache.camel.processor.idempotent.MemoryIdempotentRepository;
 import org.apache.camel.spi.IdempotentRepository;
-import org.apache.camel.util.CastUtils;
 
 /**
  * @version 
@@ -42,7 +41,7 @@ public class ManagedMemoryIdempotentCons
         MBeanServer mbeanServer = getMBeanServer();
 
         // services
-        Set<ObjectName> names = CastUtils.cast(mbeanServer.queryNames(new ObjectName("org.apache.camel" + ":type=services,*"), null));
+        Set<ObjectName> names = mbeanServer.queryNames(new ObjectName("org.apache.camel" + ":type=services,*"), null);
         ObjectName on = null;
         for (ObjectName name : names) {
             if (name.toString().contains("MemoryIdempotentRepository")) {
@@ -93,6 +92,57 @@ public class ManagedMemoryIdempotentCons
         assertTrue(repo.contains("4"));
     }
 
+    public void testDuplicateMessagesCountAreCorrectlyCounted() throws Exception {
+        MBeanServer mbeanServer = getMBeanServer();
+
+        // processors
+        Set<ObjectName> names = mbeanServer.queryNames(new ObjectName("org.apache.camel" + ":type=processors,*"), null);
+        ObjectName on = null;
+        for (ObjectName name : names) {
+            if (name.toString().contains("idempotentConsumer")) {
+                on = name;
+                break;
+            }
+        }
+        assertTrue("Should be registered", mbeanServer.isRegistered(on));
+
+        Long count = (Long) mbeanServer.getAttribute(on, "DuplicateMessageCount");
+        assertEquals(0L, count.longValue());
+        
+        resultEndpoint.expectedBodiesReceived("one", "two");
+
+        sendMessage("1", "one");
+        sendMessage("2", "two");
+        sendMessage("1", "one");
+        sendMessage("2", "two");
+
+        resultEndpoint.assertIsSatisfied();
+
+        count = (Long) mbeanServer.getAttribute(on, "DuplicateMessageCount");
+        assertEquals(2L, count.longValue());
+
+        // reset the count
+        mbeanServer.invoke(on, "resetDuplicateMessageCount", null, null);
+
+        // count should be resetted
+        count = (Long) mbeanServer.getAttribute(on, "DuplicateMessageCount");
+        assertEquals(0L, count.longValue());
+        
+        resetMocks();
+        
+        resultEndpoint.expectedBodiesReceived("five");
+
+        sendMessage("4", "four");
+        sendMessage("4", "four");
+        sendMessage("5", "five");
+        sendMessage("4", "four");
+
+        resultEndpoint.assertIsSatisfied();
+        
+        count = (Long) mbeanServer.getAttribute(on, "DuplicateMessageCount");
+        assertEquals(3L, count.longValue());
+    }
+
     protected void sendMessage(final Object messageId, final Object body) {
         template.send(startEndpoint, new Processor() {
             public void process(Exchange exchange) {