Cleanup 'handleAndReply' logic in RSocketListenerFunction
This commit is contained in:
@@ -97,16 +97,9 @@ class RSocketListenerFunction implements Function<Message<Flux<Object>>, Publish
|
|||||||
Flux<?> dataFlux =
|
Flux<?> dataFlux =
|
||||||
messageToProcess.getPayload()
|
messageToProcess.getPayload()
|
||||||
.map((payload) -> {
|
.map((payload) -> {
|
||||||
if (payload instanceof Message) {
|
return payload instanceof Message
|
||||||
return MessageBuilder.fromMessage((Message) payload).copyHeadersIfAbsent(messageToProcess.getHeaders()).build();
|
? (Message) payload
|
||||||
}
|
: MessageBuilder.withPayload(payload).copyHeadersIfAbsent(messageToProcess.getHeaders()).build();
|
||||||
else {
|
|
||||||
return MessageBuilder.withPayload(payload).copyHeadersIfAbsent(messageToProcess.getHeaders()).build();
|
|
||||||
}
|
|
||||||
// if (!(payload instanceof Message)) {
|
|
||||||
// payload = MessageBuilder.createMessage(payload, messageToProcess.getHeaders());
|
|
||||||
// }
|
|
||||||
// return payload;
|
|
||||||
});
|
});
|
||||||
if (this.targetFunction.getInputType() != null && FunctionTypeUtils.isPublisher(this.targetFunction.getInputType())) {
|
if (this.targetFunction.getInputType() != null && FunctionTypeUtils.isPublisher(this.targetFunction.getInputType())) {
|
||||||
dataFlux = dataFlux.transform((Function) this.targetFunction);
|
dataFlux = dataFlux.transform((Function) this.targetFunction);
|
||||||
|
|||||||
Reference in New Issue
Block a user