message : messages) {
- if (this.sendToChannel(message)) {
- messagesProcessed++;
- this.onSend(message);
+ boolean sent = super.sendToChannel(message);
+ if (this.source instanceof MessageDeliveryAware) {
+ if (sent) {
+ ((MessageDeliveryAware) this.source).onSend(message);
}
else {
- return messagesProcessed;
+ ((MessageDeliveryAware) this.source).onFailure(new MessageDeliveryException(message, "failed to send message"));
}
}
- return messagesProcessed;
+ return sent;
}
- /**
- * Callback method invoked after a message is sent to the channel.
- *
- * Subclasses may override. The default implementation does nothing.
- */
- protected void onSend(Message sentMessage) {
- }
-
-
- private class PollingSourceAdapterTask implements MessagingTask {
-
- public void run() {
- processMessages();
+ public void run() {
+ int messagesProcessed = 0;
+ List> messages = this.poll(this.maxMessagesPerTask);
+ for (Message> message : messages) {
+ if (this.sendMessage(message)) {
+ messagesProcessed++;
+ }
+ else {
+ break;
+ }
}
-
- public Schedule getSchedule() {
- return schedule;
+ if (logger.isDebugEnabled()) {
+ logger.debug("polling source task processed " + messagesProcessed + " messages");
}
}
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/adapter/SourceAdapter.java b/spring-integration-core/src/main/java/org/springframework/integration/adapter/SourceAdapter.java
index 5fe72ad111..58e164b3cf 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/adapter/SourceAdapter.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/adapter/SourceAdapter.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2002-2007 the original author or authors.
+ * Copyright 2002-2008 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.
@@ -16,8 +16,6 @@
package org.springframework.integration.adapter;
-import org.springframework.integration.channel.MessageChannel;
-
/**
* Base interface for source adapters.
*
@@ -25,6 +23,4 @@ import org.springframework.integration.channel.MessageChannel;
*/
public interface SourceAdapter {
- void setChannel(MessageChannel channel);
-
}
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageBus.java b/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageBus.java
index 6c096cd754..342b0e5fda 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageBus.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageBus.java
@@ -50,8 +50,8 @@ import org.springframework.integration.endpoint.MessageEndpoint;
import org.springframework.integration.handler.MessageHandler;
import org.springframework.integration.message.MessagingException;
import org.springframework.integration.scheduling.MessagePublishingErrorHandler;
+import org.springframework.integration.scheduling.MessagingTask;
import org.springframework.integration.scheduling.MessagingTaskScheduler;
-import org.springframework.integration.scheduling.MessagingTaskSchedulerAware;
import org.springframework.integration.scheduling.Schedule;
import org.springframework.integration.scheduling.SimpleMessagingTaskScheduler;
import org.springframework.integration.scheduling.Subscription;
@@ -355,8 +355,8 @@ public class MessageBus implements ChannelRegistry, EndpointRegistry, Applicatio
if (!this.initialized) {
this.initialize();
}
- if (adapter instanceof MessagingTaskSchedulerAware) {
- ((MessagingTaskSchedulerAware) adapter).setMessagingTaskScheduler(this.taskScheduler);
+ if (adapter instanceof MessagingTask) {
+ this.taskScheduler.schedule((MessagingTask) adapter);
}
if (adapter instanceof Lifecycle) {
this.lifecycleSourceAdapters.add((Lifecycle) adapter);
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/ChannelAdapterParser.java b/spring-integration-core/src/main/java/org/springframework/integration/config/ChannelAdapterParser.java
index 9ccaac21b6..a2f32b0e5b 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/config/ChannelAdapterParser.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/config/ChannelAdapterParser.java
@@ -30,6 +30,7 @@ import org.springframework.integration.adapter.MethodInvokingSource;
import org.springframework.integration.adapter.MethodInvokingTarget;
import org.springframework.integration.adapter.PollingSourceAdapter;
import org.springframework.integration.endpoint.DefaultMessageEndpoint;
+import org.springframework.integration.scheduling.PollingSchedule;
import org.springframework.integration.scheduling.Subscription;
import org.springframework.util.StringUtils;
@@ -77,21 +78,22 @@ public class ChannelAdapterParser implements BeanDefinitionParser {
if (this.isInbound) {
adapterDef = new RootBeanDefinition(PollingSourceAdapter.class);
invokerDef = new RootBeanDefinition(MethodInvokingSource.class);
+ String invokerBeanName = this.configureAndRegisterInvoker(invokerDef, ref, method, parserContext);
String period = element.getAttribute(PERIOD_ATTRIBUTE);
- if (StringUtils.hasText(period)) {
- adapterDef.getPropertyValues().addPropertyValue("period", period);
+ if (!StringUtils.hasText(period)) {
+ throw new ConfigurationException("'period' is required");
}
- adapterDef.getPropertyValues().addPropertyValue("channel", new RuntimeBeanReference(channel));
+ PollingSchedule schedule = new PollingSchedule(Integer.valueOf(period));
+ adapterDef.getConstructorArgumentValues().addGenericArgumentValue(new RuntimeBeanReference(invokerBeanName));
+ adapterDef.getConstructorArgumentValues().addGenericArgumentValue(new RuntimeBeanReference(channel));
+ adapterDef.getConstructorArgumentValues().addGenericArgumentValue(schedule);
}
else {
adapterDef = new RootBeanDefinition(DefaultTargetAdapter.class);
invokerDef = new RootBeanDefinition(MethodInvokingTarget.class);
+ String invokerBeanName = this.configureAndRegisterInvoker(invokerDef, ref, method, parserContext);
+ adapterDef.getConstructorArgumentValues().addGenericArgumentValue(new RuntimeBeanReference(invokerBeanName));
}
- invokerDef.getPropertyValues().addPropertyValue("object", new RuntimeBeanReference(ref));
- invokerDef.getPropertyValues().addPropertyValue("method", method);
- String invokerBeanName = parserContext.getReaderContext().generateBeanName(invokerDef);
- parserContext.registerBeanComponent(new BeanComponentDefinition(invokerDef, invokerBeanName));
- adapterDef.getConstructorArgumentValues().addGenericArgumentValue(new RuntimeBeanReference(invokerBeanName));
adapterDef.setSource(parserContext.extractSource(element));
String beanName = element.getAttribute(ID_ATTRIBUTE);
if (!StringUtils.hasText(beanName)) {
@@ -112,4 +114,12 @@ public class ChannelAdapterParser implements BeanDefinitionParser {
return adapterDef;
}
+ private String configureAndRegisterInvoker(RootBeanDefinition invokerDef, String objectRef, String methodName, ParserContext parserContext) {
+ invokerDef.getPropertyValues().addPropertyValue("object", new RuntimeBeanReference(objectRef));
+ invokerDef.getPropertyValues().addPropertyValue("method", methodName);
+ String invokerBeanName = parserContext.getReaderContext().generateBeanName(invokerDef);
+ parserContext.registerBeanComponent(new BeanComponentDefinition(invokerDef, invokerBeanName));
+ return invokerBeanName;
+ }
+
}
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/MessageEndpointAnnotationPostProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/config/MessageEndpointAnnotationPostProcessor.java
index 80ef182b8b..c0ed1a23c6 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/config/MessageEndpointAnnotationPostProcessor.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/config/MessageEndpointAnnotationPostProcessor.java
@@ -49,8 +49,7 @@ import org.springframework.integration.annotation.Router;
import org.springframework.integration.annotation.Splitter;
import org.springframework.integration.bus.MessageBus;
import org.springframework.integration.channel.ChannelRegistryAware;
-import org.springframework.integration.channel.MessageChannel;
-import org.springframework.integration.channel.SimpleChannel;
+import org.springframework.integration.dispatcher.SynchronousChannel;
import org.springframework.integration.endpoint.ConcurrencyPolicy;
import org.springframework.integration.endpoint.DefaultMessageEndpoint;
import org.springframework.integration.handler.AbstractMessageHandlerAdapter;
@@ -161,17 +160,15 @@ public class MessageEndpointAnnotationPostProcessor implements BeanPostProcessor
MethodInvokingSource