GH-1835 Fixed channel-to-publisher adapter
- Needed to copy one from SI to address the issue but the actual fix should go into SI Resolves #1835
This commit is contained in:
@@ -34,7 +34,9 @@ import java.util.stream.Stream;
|
||||
import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
import org.reactivestreams.Publisher;
|
||||
import org.reactivestreams.Subscriber;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.FluxSink;
|
||||
import reactor.core.publisher.Mono;
|
||||
import reactor.core.publisher.MonoSink;
|
||||
import reactor.util.function.Tuples;
|
||||
@@ -82,7 +84,6 @@ import org.springframework.core.annotation.AnnotationUtils;
|
||||
import org.springframework.core.env.Environment;
|
||||
import org.springframework.core.type.MethodMetadata;
|
||||
import org.springframework.integration.channel.AbstractMessageChannel;
|
||||
import org.springframework.integration.channel.MessageChannelReactiveUtils;
|
||||
import org.springframework.integration.dsl.IntegrationFlow;
|
||||
import org.springframework.integration.dsl.IntegrationFlowBuilder;
|
||||
import org.springframework.integration.dsl.IntegrationFlows;
|
||||
@@ -91,6 +92,7 @@ import org.springframework.integration.support.MessageBuilder;
|
||||
import org.springframework.lang.Nullable;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.MessageChannel;
|
||||
import org.springframework.messaging.MessageHandler;
|
||||
import org.springframework.messaging.MessageHeaders;
|
||||
import org.springframework.messaging.SubscribableChannel;
|
||||
import org.springframework.util.Assert;
|
||||
@@ -364,7 +366,7 @@ public class FunctionConfiguration {
|
||||
|
||||
Publisher[] inputPublishers = inputBindingNames.stream().map(inputBindingName -> {
|
||||
SubscribableChannel inputChannel = this.context.getBean(inputBindingName, SubscribableChannel.class);
|
||||
return this.enhancePublisher(MessageChannelReactiveUtils.toPublisher(inputChannel), inputBindingName);
|
||||
return this.enhancePublisher(new SubscribableChannelPublisherAdapter(inputChannel), inputBindingName);
|
||||
}).toArray(Publisher[]::new);
|
||||
|
||||
|
||||
@@ -425,7 +427,7 @@ public class FunctionConfiguration {
|
||||
if (FunctionTypeUtils.isReactive(FunctionTypeUtils.getInputType(functionType, 0)) && StringUtils.hasText(outputChannelName)) {
|
||||
MessageChannel outputChannel = context.getBean(outputChannelName, MessageChannel.class);
|
||||
SubscribableChannel subscribeChannel = (SubscribableChannel) inputChannel;
|
||||
Publisher<?> publisher = this.enhancePublisher(MessageChannelReactiveUtils.toPublisher(subscribeChannel),
|
||||
Publisher<?> publisher = this.enhancePublisher(new SubscribableChannelPublisherAdapter<>(subscribeChannel),
|
||||
((DirectWithAttributesChannel) inputChannel).getBeanName());
|
||||
this.subscribeToInput(function, publisher, message -> outputChannel.send((Message<?>) message));
|
||||
}
|
||||
@@ -677,4 +679,27 @@ public class FunctionConfiguration {
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
private static final class SubscribableChannelPublisherAdapter<T> implements Publisher<Message<T>> {
|
||||
|
||||
private final SubscribableChannel channel;
|
||||
|
||||
SubscribableChannelPublisherAdapter(SubscribableChannel channel) {
|
||||
this.channel = channel;
|
||||
}
|
||||
|
||||
@Override
|
||||
@SuppressWarnings("unchecked")
|
||||
public void subscribe(Subscriber<? super Message<T>> subscriber) {
|
||||
Flux.
|
||||
<Message<?>>push(emitter -> {
|
||||
MessageHandler messageHandler = emitter::next;
|
||||
this.channel.subscribe(messageHandler);
|
||||
emitter.onCancel(() -> this.channel.unsubscribe(messageHandler));
|
||||
},
|
||||
FluxSink.OverflowStrategy.BUFFER)
|
||||
.subscribe((Subscriber<? super Message<?>>) subscriber);
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user