From 1f22d47b5fa91fbe2c5d0d79913f054ec6b3b26a Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Mon, 1 Feb 2010 21:25:09 +0000 Subject: [PATCH] initial commit for NotificationListeningAdapter --- .../integration/jmx/JmxHeaders.java | 31 ++++ .../jmx/NotificationListeningAdapter.java | 115 ++++++++++++ .../NotificationListeningAdapterTests.java | 167 ++++++++++++++++++ 3 files changed, 313 insertions(+) create mode 100644 org.springframework.integration.jmx/src/main/java/org/springframework/integration/jmx/JmxHeaders.java create mode 100644 org.springframework.integration.jmx/src/main/java/org/springframework/integration/jmx/NotificationListeningAdapter.java create mode 100644 org.springframework.integration.jmx/src/test/java/org/springframework/integration/jmx/NotificationListeningAdapterTests.java diff --git a/org.springframework.integration.jmx/src/main/java/org/springframework/integration/jmx/JmxHeaders.java b/org.springframework.integration.jmx/src/main/java/org/springframework/integration/jmx/JmxHeaders.java new file mode 100644 index 0000000000..017ea36a06 --- /dev/null +++ b/org.springframework.integration.jmx/src/main/java/org/springframework/integration/jmx/JmxHeaders.java @@ -0,0 +1,31 @@ +/* + * 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 org.springframework.integration.core.MessageHeaders; + +/** + * @author Mark Fisher + * @since 2.0 + */ +public abstract class JmxHeaders { + + private static final String PREFIX = MessageHeaders.PREFIX + "jmx"; + + public static final String NOTIFICATION_HANDBACK = PREFIX + "_handback"; + +} diff --git a/org.springframework.integration.jmx/src/main/java/org/springframework/integration/jmx/NotificationListeningAdapter.java b/org.springframework.integration.jmx/src/main/java/org/springframework/integration/jmx/NotificationListeningAdapter.java new file mode 100644 index 0000000000..cf8bc02d30 --- /dev/null +++ b/org.springframework.integration.jmx/src/main/java/org/springframework/integration/jmx/NotificationListeningAdapter.java @@ -0,0 +1,115 @@ +/* + * 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 java.util.Arrays; +import java.util.LinkedHashSet; +import java.util.Set; + +import javax.management.InstanceNotFoundException; +import javax.management.ListenerNotFoundException; +import javax.management.MBeanServer; +import javax.management.Notification; +import javax.management.NotificationFilter; +import javax.management.NotificationListener; +import javax.management.ObjectName; + +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; + +import org.springframework.integration.core.Message; +import org.springframework.integration.core.MessageHistory.ComponentType; +import org.springframework.integration.endpoint.MessageProducerSupport; +import org.springframework.integration.message.MessageBuilder; +import org.springframework.util.Assert; + +/** + * @author Mark Fisher + * @since 2.0 + */ +public class NotificationListeningAdapter extends MessageProducerSupport implements NotificationListener { + + private final Log logger = LogFactory.getLog(this.getClass()); + + private volatile MBeanServer server; + + private volatile Set objectNames; + + private volatile NotificationFilter filter; + + private volatile Object handback; + + + public void setServer(MBeanServer server) { + this.server = server; + } + + public void setObjectNames(ObjectName... objectNames) { + this.objectNames = new LinkedHashSet(Arrays.asList(objectNames)); + } + + public void setFilter(NotificationFilter filter) { + this.filter = filter; + } + + public void setHandback(Object handback) { + this.handback = handback; + } + + public void handleNotification(Notification notification, Object handback) { + if (logger.isInfoEnabled()) { + logger.info("received notification: " + notification + ", and handback: " + handback); + } + MessageBuilder builder = MessageBuilder.withPayload(notification); + if (handback != null) { + builder.setHeader(JmxHeaders.NOTIFICATION_HANDBACK, handback); + } + Message message = builder.build(); + message.getHeaders().getHistory().add(ComponentType.endpoint, this.getBeanName()); + this.sendMessage(message); + } + + @Override + protected void doStart() { + try { + Assert.notNull(this.server, "MBeanServer is required."); + for (ObjectName objectName : this.objectNames) { + this.server.addNotificationListener(objectName, this, this.filter, this.handback); + } + } + catch (InstanceNotFoundException e) { + throw new IllegalStateException("Failed to find MBean instance.", e); + } + } + + @Override + protected void doStop() { + try { + Assert.notNull(this.server, "MBeanServer is required."); + for (ObjectName objectName : this.objectNames) { + this.server.removeNotificationListener(objectName, this, this.filter, this.handback); + } + } + catch (InstanceNotFoundException e) { + throw new IllegalStateException("Failed to find MBean instance.", e); + } + catch (ListenerNotFoundException e) { + throw new IllegalStateException("Failed to find NotificationListener.", e); + } + } + +} diff --git a/org.springframework.integration.jmx/src/test/java/org/springframework/integration/jmx/NotificationListeningAdapterTests.java b/org.springframework.integration.jmx/src/test/java/org/springframework/integration/jmx/NotificationListeningAdapterTests.java new file mode 100644 index 0000000000..a938ad02bd --- /dev/null +++ b/org.springframework.integration.jmx/src/test/java/org/springframework/integration/jmx/NotificationListeningAdapterTests.java @@ -0,0 +1,167 @@ +/* + * 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 static org.junit.Assert.assertNull; +import static org.junit.Assert.assertTrue; + +import java.util.concurrent.atomic.AtomicInteger; + +import javax.management.MBeanServer; +import javax.management.Notification; +import javax.management.NotificationFilter; +import javax.management.ObjectName; + +import org.junit.After; +import org.junit.Before; +import org.junit.Test; + +import org.springframework.integration.channel.QueueChannel; +import org.springframework.integration.core.Message; +import org.springframework.jmx.export.MBeanExporter; +import org.springframework.jmx.export.notification.NotificationPublisher; +import org.springframework.jmx.export.notification.NotificationPublisherAware; +import org.springframework.jmx.support.MBeanServerFactoryBean; +import org.springframework.jmx.support.ObjectNameManager; + +/** + * @author Mark Fisher + * @since 2.0 + */ +public class NotificationListeningAdapterTests { + + private volatile MBeanServer server; + + private volatile ObjectName objectName; + + private final NumberHolder numberHolder = new NumberHolder(); + + + @Before + public void setup() throws Exception { + MBeanServerFactoryBean serverFactoryBean = new MBeanServerFactoryBean(); + serverFactoryBean.setLocateExistingServerIfPossible(true); + serverFactoryBean.afterPropertiesSet(); + this.server = serverFactoryBean.getObject(); + MBeanExporter exporter = new MBeanExporter(); + exporter.setAutodetect(false); + exporter.afterPropertiesSet(); + this.objectName = ObjectNameManager.getInstance("si:name=numberHolder"); + exporter.registerManagedResource(this.numberHolder, this.objectName); + } + + @After + public void cleanup() throws Exception { + this.server.unregisterMBean(this.objectName); + } + + @Test + public void simpleNotification() throws Exception { + QueueChannel outputChannel = new QueueChannel(); + NotificationListeningAdapter adapter = new NotificationListeningAdapter(); + adapter.setServer(this.server); + adapter.setObjectNames(this.objectName); + adapter.setOutputChannel(outputChannel); + adapter.afterPropertiesSet(); + adapter.start(); + this.numberHolder.publish("foo"); + Message message = outputChannel.receive(0); + assertNotNull(message); + assertTrue(message.getPayload() instanceof Notification); + Notification notification = (Notification) message.getPayload(); + assertEquals("foo", notification.getMessage()); + assertEquals(objectName, notification.getSource()); + assertNull(message.getHeaders().get(JmxHeaders.NOTIFICATION_HANDBACK)); + } + + @Test + public void notificationWithHandback() throws Exception { + QueueChannel outputChannel = new QueueChannel(); + NotificationListeningAdapter adapter = new NotificationListeningAdapter(); + adapter.setServer(this.server); + adapter.setObjectNames(this.objectName); + adapter.setOutputChannel(outputChannel); + Integer handback = new Integer(123); + adapter.setHandback(handback); + adapter.afterPropertiesSet(); + adapter.start(); + this.numberHolder.publish("foo"); + Message message = outputChannel.receive(0); + assertNotNull(message); + assertTrue(message.getPayload() instanceof Notification); + Notification notification = (Notification) message.getPayload(); + assertEquals("foo", notification.getMessage()); + assertEquals(objectName, notification.getSource()); + assertEquals(handback, message.getHeaders().get(JmxHeaders.NOTIFICATION_HANDBACK)); + } + + @Test + @SuppressWarnings("serial") + public void notificationWithFilter() throws Exception { + QueueChannel outputChannel = new QueueChannel(); + NotificationListeningAdapter adapter = new NotificationListeningAdapter(); + adapter.setServer(this.server); + adapter.setObjectNames(this.objectName); + adapter.setOutputChannel(outputChannel); + adapter.setFilter(new NotificationFilter() { + public boolean isNotificationEnabled(Notification notification) { + return !notification.getMessage().equals("bad"); + } + }); + adapter.afterPropertiesSet(); + adapter.start(); + this.numberHolder.publish("bad"); + Message message = outputChannel.receive(0); + assertNull(message); + this.numberHolder.publish("okay"); + message = outputChannel.receive(0); + assertNotNull(message); + assertTrue(message.getPayload() instanceof Notification); + Notification notification = (Notification) message.getPayload(); + assertEquals("okay", notification.getMessage()); + } + + + public static class NumberHolder implements NotificationPublisherAware { + + private final AtomicInteger number = new AtomicInteger(); + + private final AtomicInteger sequence = new AtomicInteger(); + + private volatile NotificationPublisher notificationPublisher; + + public int getNumber() { + return this.number.get(); + } + + public void setNumber(int value) { + this.number.set(value); + } + + public void setNotificationPublisher(NotificationPublisher notificationPublisher) { + this.notificationPublisher = notificationPublisher; + } + + public void publish(String message) { + Notification notification = new Notification("testType", this, sequence.getAndIncrement(), message); + this.notificationPublisher.sendNotification(notification); + } + } + +}