From 959a8614a70df19de776a10aa2e92fbb583e8491 Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Tue, 2 Feb 2010 17:44:33 +0000 Subject: [PATCH] INT-978 initial commit of AttributePollingMessageSource --- .../jmx/AttributePollingMessageSource.java | 74 ++++++++++++++++ .../AttributePollingMessageSourceTests.java | 88 +++++++++++++++++++ 2 files changed, 162 insertions(+) create mode 100644 org.springframework.integration.jmx/src/main/java/org/springframework/integration/jmx/AttributePollingMessageSource.java create mode 100644 org.springframework.integration.jmx/src/test/java/org/springframework/integration/jmx/AttributePollingMessageSourceTests.java diff --git a/org.springframework.integration.jmx/src/main/java/org/springframework/integration/jmx/AttributePollingMessageSource.java b/org.springframework.integration.jmx/src/main/java/org/springframework/integration/jmx/AttributePollingMessageSource.java new file mode 100644 index 0000000000..0abf9a3d27 --- /dev/null +++ b/org.springframework.integration.jmx/src/main/java/org/springframework/integration/jmx/AttributePollingMessageSource.java @@ -0,0 +1,74 @@ +/* + * Copyright 2002-2010 the original author or authors. + * + * Licensed 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.springframework.integration.jmx; + +import javax.management.MBeanServer; +import javax.management.MalformedObjectNameException; +import javax.management.ObjectName; + +import org.springframework.integration.core.Message; +import org.springframework.integration.core.MessagingException; +import org.springframework.integration.gateway.SimpleMessageMapper; +import org.springframework.integration.message.InboundMessageMapper; +import org.springframework.integration.message.MessageSource; +import org.springframework.jmx.support.ObjectNameManager; + +/** + * @author Mark Fisher + * @since 2.0 + */ +@SuppressWarnings("unchecked") +public class AttributePollingMessageSource implements MessageSource { + + private volatile ObjectName objectName; + + private volatile String attributeName; + + private volatile MBeanServer server; + + private volatile InboundMessageMapper mapper = new SimpleMessageMapper(); + + + public void setObjectName(String objectName) { + try { + this.objectName = ObjectNameManager.getInstance(objectName); + } + catch (MalformedObjectNameException e) { + throw new IllegalArgumentException(e); + } + } + + public void setAttributeName(String attributeName) { + this.attributeName = attributeName; + } + + public void setServer(MBeanServer server) { + this.server = server; + } + + public Message receive() { + try { + Object value = this.server.getAttribute(this.objectName, this.attributeName); + return this.mapper.toMessage(value); + } + catch (Exception e) { + throw new MessagingException("failed to retrieve JMX attribute '" + + this.attributeName + "' on MBean [" + this.objectName + "]", e); + } + } + +} diff --git a/org.springframework.integration.jmx/src/test/java/org/springframework/integration/jmx/AttributePollingMessageSourceTests.java b/org.springframework.integration.jmx/src/test/java/org/springframework/integration/jmx/AttributePollingMessageSourceTests.java new file mode 100644 index 0000000000..007414b2e4 --- /dev/null +++ b/org.springframework.integration.jmx/src/test/java/org/springframework/integration/jmx/AttributePollingMessageSourceTests.java @@ -0,0 +1,88 @@ +/* + * Copyright 2002-2010 the original author or authors. + * + * Licensed 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.springframework.integration.jmx; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; + +import java.util.concurrent.atomic.AtomicInteger; + +import javax.management.MBeanServer; + +import org.junit.Before; +import org.junit.Test; + +import org.springframework.integration.core.Message; +import org.springframework.jmx.support.MBeanServerFactoryBean; +import org.springframework.jmx.support.ObjectNameManager; + +/** + * @author Mark Fisher + * @since 2.0 + */ +public class AttributePollingMessageSourceTests { + + private final TestCounter counter = new TestCounter(); + + private volatile MBeanServer server; + + + @Before + public void setup() throws Exception { + MBeanServerFactoryBean factoryBean = new MBeanServerFactoryBean(); + factoryBean.setLocateExistingServerIfPossible(true); + factoryBean.afterPropertiesSet(); + this.server = factoryBean.getObject(); + this.server.registerMBean(this.counter, ObjectNameManager.getInstance("test:name=counter")); + } + + + @Test + public void basicPolling() { + AttributePollingMessageSource source = new AttributePollingMessageSource(); + source.setAttributeName("Count"); + source.setObjectName("test:name=counter"); + source.setServer(this.server); + Message message1 = source.receive(); + assertNotNull(message1); + assertEquals(0, message1.getPayload()); + this.counter.increment(); + Message message2 = source.receive(); + assertNotNull(message2); + assertEquals(1, message2.getPayload()); + } + + + public static interface TestCounterMBean { + int getCount(); + } + + + public static class TestCounter implements TestCounterMBean { + + private final AtomicInteger counter = new AtomicInteger(); + + public int getCount() { + return this.counter.get(); + } + + public void increment() { + this.counter.incrementAndGet(); + } + } + +}