Author: bvahdat
Date: Sun Feb 26 15:44:57 2012
New Revision: 1293856
URL: http://svn.apache.org/viewvc?rev=1293856&view=rev
Log:
Merged revisions 1293852 via svnmerge from
https://svn.apache.org/repos/asf/camel/trunk
........
r1293852 | bvahdat | 2012-02-26 16:29:59 +0100 (So, 26 Feb 2012) | 1 line
CAMEL-4782: Add ManagedIdempotentConsumerBean providing get & reset of the
duplicate messages count detected.
........
Added:
camel/branches/camel-2.9.x/camel-core/src/main/java/org/apache/camel/api/management/mbean/ManagedIdempotentConsumerMBean.java
- copied unchanged from r1293852,
camel/trunk/camel-core/src/main/java/org/apache/camel/api/management/mbean/ManagedIdempotentConsumerMBean.java
camel/branches/camel-2.9.x/camel-core/src/main/java/org/apache/camel/management/mbean/ManagedIdempotentConsumer.java
- copied unchanged from r1293852,
camel/trunk/camel-core/src/main/java/org/apache/camel/management/mbean/ManagedIdempotentConsumer.java
Modified:
camel/branches/camel-2.9.x/ (props changed)
camel/branches/camel-2.9.x/camel-core/src/main/java/org/apache/camel/management/DefaultManagementLifecycleStrategy.java
camel/branches/camel-2.9.x/camel-core/src/main/java/org/apache/camel/management/DefaultManagementObjectStrategy.java
camel/branches/camel-2.9.x/camel-core/src/main/java/org/apache/camel/processor/idempotent/IdempotentConsumer.java
camel/branches/camel-2.9.x/camel-core/src/test/java/org/apache/camel/management/ManagedMemoryIdempotentConsumerTest.java
Propchange: camel/branches/camel-2.9.x/
------------------------------------------------------------------------------
--- svn:mergeinfo (original)
+++ svn:mergeinfo Sun Feb 26 15:44:57 2012
@@ -1 +1 @@
-/camel/trunk:1243046,1243057,1243234,1244518,1244644,1244859,1244861,1244864,1244870,1244872,1245021,1291555,1291727,1291848,1291864,1292114,1292384,1292725,1292760,1292767,1293079,1293268,1293288,1293330,1293590
+/camel/trunk:1243046,1243057,1243234,1244518,1244644,1244859,1244861,1244864,1244870,1244872,1245021,1291555,1291727,1291848,1291864,1292114,1292384,1292725,1292760,1292767,1293079,1293268,1293288,1293330,1293590,1293852
Propchange: camel/branches/camel-2.9.x/
------------------------------------------------------------------------------
Binary property 'svnmerge-integrated' - no diff available.
Modified:
camel/branches/camel-2.9.x/camel-core/src/main/java/org/apache/camel/management/DefaultManagementLifecycleStrategy.java
URL:
http://svn.apache.org/viewvc/camel/branches/camel-2.9.x/camel-core/src/main/java/org/apache/camel/management/DefaultManagementLifecycleStrategy.java?rev=1293856&r1=1293855&r2=1293856&view=diff
==============================================================================
---
camel/branches/camel-2.9.x/camel-core/src/main/java/org/apache/camel/management/DefaultManagementLifecycleStrategy.java
(original)
+++
camel/branches/camel-2.9.x/camel-core/src/main/java/org/apache/camel/management/DefaultManagementLifecycleStrategy.java
Sun Feb 26 15:44:57 2012
@@ -430,10 +430,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/branches/camel-2.9.x/camel-core/src/main/java/org/apache/camel/management/DefaultManagementObjectStrategy.java
URL:
http://svn.apache.org/viewvc/camel/branches/camel-2.9.x/camel-core/src/main/java/org/apache/camel/management/DefaultManagementObjectStrategy.java?rev=1293856&r1=1293855&r2=1293856&view=diff
==============================================================================
---
camel/branches/camel-2.9.x/camel-core/src/main/java/org/apache/camel/management/DefaultManagementObjectStrategy.java
(original)
+++
camel/branches/camel-2.9.x/camel-core/src/main/java/org/apache/camel/management/DefaultManagementObjectStrategy.java
Sun Feb 26 15:44:57 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);
}
Modified:
camel/branches/camel-2.9.x/camel-core/src/main/java/org/apache/camel/processor/idempotent/IdempotentConsumer.java
URL:
http://svn.apache.org/viewvc/camel/branches/camel-2.9.x/camel-core/src/main/java/org/apache/camel/processor/idempotent/IdempotentConsumer.java?rev=1293856&r1=1293855&r2=1293856&view=diff
==============================================================================
---
camel/branches/camel-2.9.x/camel-core/src/main/java/org/apache/camel/processor/idempotent/IdempotentConsumer.java
(original)
+++
camel/branches/camel-2.9.x/camel-core/src/main/java/org/apache/camel/processor/idempotent/IdempotentConsumer.java
Sun Feb 26 15:44:57 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/branches/camel-2.9.x/camel-core/src/test/java/org/apache/camel/management/ManagedMemoryIdempotentConsumerTest.java
URL:
http://svn.apache.org/viewvc/camel/branches/camel-2.9.x/camel-core/src/test/java/org/apache/camel/management/ManagedMemoryIdempotentConsumerTest.java?rev=1293856&r1=1293855&r2=1293856&view=diff
==============================================================================
---
camel/branches/camel-2.9.x/camel-core/src/test/java/org/apache/camel/management/ManagedMemoryIdempotentConsumerTest.java
(original)
+++
camel/branches/camel-2.9.x/camel-core/src/test/java/org/apache/camel/management/ManagedMemoryIdempotentConsumerTest.java
Sun Feb 26 15:44:57 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) {