More Kotlin DSL improvements
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2019-2021 the original author or authors.
|
||||
* Copyright 2019-2022 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.
|
||||
@@ -496,7 +496,7 @@ public abstract class BaseIntegrationFlowDefinition<B extends BaseIntegrationFlo
|
||||
* <pre class="code">
|
||||
* {@code
|
||||
* .transform("payload")
|
||||
* .wireTap(new WireTap(tapChannel().selector(m -> m.getPayload().equals("foo")))
|
||||
* .wireTap(new WireTap(tapChannel()).selector(m -> m.getPayload().equals("foo")))
|
||||
* .channel("foo")
|
||||
* }
|
||||
* </pre>
|
||||
@@ -1721,7 +1721,7 @@ public abstract class BaseIntegrationFlowDefinition<B extends BaseIntegrationFlo
|
||||
* Populate the
|
||||
* {@link org.springframework.integration.aggregator.ResequencingMessageHandler} with
|
||||
* provided options from {@link ResequencerSpec}.
|
||||
* In addition accept options for the integration endpoint using {@link GenericEndpointSpec}.
|
||||
* In addition, accept options for the integration endpoint using {@link GenericEndpointSpec}.
|
||||
* Typically used with a Java 8 Lambda expression:
|
||||
* <pre class="code">
|
||||
* {@code
|
||||
@@ -1744,7 +1744,7 @@ public abstract class BaseIntegrationFlowDefinition<B extends BaseIntegrationFlo
|
||||
* @return the current {@link BaseIntegrationFlowDefinition}.
|
||||
*/
|
||||
public B aggregate() {
|
||||
return aggregate(null);
|
||||
return aggregate((Consumer<AggregatorSpec>) null);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1760,7 +1760,7 @@ public abstract class BaseIntegrationFlowDefinition<B extends BaseIntegrationFlo
|
||||
|
||||
/**
|
||||
* Populate the {@link AggregatingMessageHandler} with provided options from {@link AggregatorSpec}.
|
||||
* In addition accept options for the integration endpoint using {@link GenericEndpointSpec}.
|
||||
* In addition, accept options for the integration endpoint using {@link GenericEndpointSpec}.
|
||||
* Typically used with a Java 8 Lambda expression:
|
||||
* <pre class="code">
|
||||
* {@code
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2020-2021 the original author or authors.
|
||||
* Copyright 2020-2022 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.
|
||||
@@ -22,7 +22,6 @@ import org.springframework.integration.endpoint.MessageProducerSupport
|
||||
import org.springframework.integration.gateway.MessagingGatewaySupport
|
||||
import org.springframework.messaging.Message
|
||||
import org.springframework.messaging.MessageChannel
|
||||
import java.util.function.Consumer
|
||||
|
||||
private fun buildIntegrationFlow(flowBuilder: IntegrationFlowBuilder,
|
||||
flow: (KotlinIntegrationFlowDefinition) -> Unit) =
|
||||
@@ -160,6 +159,8 @@ fun integrationFlow(producerSpec: MessageProducerSpec<*, *>,
|
||||
* `IntegrationFlows.from(IntegrationFlow)` factory method.
|
||||
*
|
||||
* @author Artem Bilan
|
||||
*
|
||||
* @since 5.5.8
|
||||
*/
|
||||
fun integrationFlow(sourceFlow: IntegrationFlow, flow: KotlinIntegrationFlowDefinition.() -> Unit) =
|
||||
buildIntegrationFlow(IntegrationFlows.from(sourceFlow), flow)
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2020-2021 the original author or authors.
|
||||
* Copyright 2020-2022 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.
|
||||
@@ -55,6 +55,7 @@ import org.springframework.messaging.Message
|
||||
import org.springframework.messaging.MessageChannel
|
||||
import org.springframework.messaging.MessageHandler
|
||||
import org.springframework.messaging.MessageHeaders
|
||||
import org.springframework.messaging.support.ChannelInterceptor
|
||||
import reactor.core.publisher.Flux
|
||||
import java.util.function.Consumer
|
||||
|
||||
@@ -75,7 +76,8 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* with reified generic type.
|
||||
*/
|
||||
inline fun <reified T> convert(
|
||||
crossinline configurer: GenericEndpointSpec<MessageTransformingHandler>.() -> Unit = {}) {
|
||||
crossinline configurer: GenericEndpointSpec<MessageTransformingHandler>.() -> Unit = {}
|
||||
) {
|
||||
|
||||
this.delegate.convert(T::class.java) { configurer(it) }
|
||||
}
|
||||
@@ -93,8 +95,9 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* with reified generic type.
|
||||
*/
|
||||
inline fun <reified P> transform(
|
||||
crossinline function: (P) -> Any,
|
||||
crossinline configurer: GenericEndpointSpec<MessageTransformingHandler>.() -> Unit) {
|
||||
crossinline function: (P) -> Any,
|
||||
crossinline configurer: GenericEndpointSpec<MessageTransformingHandler>.() -> Unit
|
||||
) {
|
||||
|
||||
this.delegate.transform(P::class.java, { function(it) }) { configurer(it) }
|
||||
}
|
||||
@@ -113,8 +116,9 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* with reified generic type.
|
||||
*/
|
||||
inline fun <reified P> split(
|
||||
crossinline function: (P) -> Any,
|
||||
crossinline configurer: KotlinSplitterEndpointSpec<MethodInvokingSplitter>.() -> Unit) {
|
||||
crossinline function: (P) -> Any,
|
||||
crossinline configurer: KotlinSplitterEndpointSpec<MethodInvokingSplitter>.() -> Unit
|
||||
) {
|
||||
|
||||
this.delegate.split(P::class.java, { function(it) }) { configurer(KotlinSplitterEndpointSpec(it)) }
|
||||
}
|
||||
@@ -132,8 +136,9 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* with reified generic type.
|
||||
*/
|
||||
inline fun <reified P> filter(
|
||||
crossinline function: (P) -> Boolean,
|
||||
crossinline filterConfigurer: KotlinFilterEndpointSpec.() -> Unit) {
|
||||
crossinline function: (P) -> Boolean,
|
||||
crossinline filterConfigurer: KotlinFilterEndpointSpec.() -> Unit
|
||||
) {
|
||||
|
||||
this.delegate.filter(P::class.java, { function(it) }) { filterConfigurer(KotlinFilterEndpointSpec(it)) }
|
||||
}
|
||||
@@ -152,8 +157,9 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* with reified generic type.
|
||||
*/
|
||||
inline fun <reified P, T> route(
|
||||
crossinline function: (P) -> T,
|
||||
crossinline configurer: KotlinRouterSpec<T, MethodInvokingRouter>.() -> Unit) {
|
||||
crossinline function: (P) -> T,
|
||||
crossinline configurer: KotlinRouterSpec<T, MethodInvokingRouter>.() -> Unit
|
||||
) {
|
||||
|
||||
this.delegate.route(P::class.java, { function(it) }) { configurer(KotlinRouterSpec(it)) }
|
||||
}
|
||||
@@ -212,15 +218,17 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* The [org.springframework.integration.channel.BroadcastCapableChannel] `channel()`
|
||||
* method specific implementation to allow the use of the 'subflow' subscriber capability.
|
||||
*/
|
||||
fun publishSubscribe(broadcastCapableChannel: BroadcastCapableChannel,
|
||||
vararg subscribeSubFlows: KotlinIntegrationFlowDefinition.() -> Unit) {
|
||||
fun publishSubscribe(
|
||||
broadcastCapableChannel: BroadcastCapableChannel,
|
||||
vararg subscribeSubFlows: KotlinIntegrationFlowDefinition.() -> Unit
|
||||
) {
|
||||
|
||||
val publishSubscribeChannelConfigurer =
|
||||
Consumer<BroadcastPublishSubscribeSpec> { spec ->
|
||||
subscribeSubFlows.forEach { subFlow ->
|
||||
spec.subscribe { subFlow(KotlinIntegrationFlowDefinition(it)) }
|
||||
}
|
||||
Consumer<BroadcastPublishSubscribeSpec> { spec ->
|
||||
subscribeSubFlows.forEach { subFlow ->
|
||||
spec.subscribe { subFlow(KotlinIntegrationFlowDefinition(it)) }
|
||||
}
|
||||
}
|
||||
this.delegate.publishSubscribeChannel(broadcastCapableChannel, publishSubscribeChannelConfigurer)
|
||||
}
|
||||
|
||||
@@ -244,8 +252,9 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
*/
|
||||
fun wireTap(wireTapConfigurer: WireTapSpec.() -> Unit, flow: KotlinIntegrationFlowDefinition.() -> Unit) {
|
||||
this.delegate.wireTap(
|
||||
IntegrationFlow { flow(KotlinIntegrationFlowDefinition(it)) },
|
||||
Consumer(wireTapConfigurer))
|
||||
IntegrationFlow { flow(KotlinIntegrationFlowDefinition(it)) },
|
||||
Consumer(wireTapConfigurer)
|
||||
)
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -295,8 +304,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* for the provided `Transformer` instance.
|
||||
* @since 5.3.1
|
||||
*/
|
||||
fun transform(transformer: Transformer,
|
||||
endpointConfigurer: GenericEndpointSpec<MessageTransformingHandler>.() -> Unit = {}) {
|
||||
fun transform(
|
||||
transformer: Transformer,
|
||||
endpointConfigurer: GenericEndpointSpec<MessageTransformingHandler>.() -> Unit = {}
|
||||
) {
|
||||
|
||||
this.delegate.transform(transformer) { endpointConfigurer(it) }
|
||||
}
|
||||
@@ -305,8 +316,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* Populate the [Transformer] EI Pattern specific [MessageHandler] implementation
|
||||
* for the SpEL [Expression].
|
||||
*/
|
||||
fun transform(expression: String,
|
||||
endpointConfigurer: GenericEndpointSpec<MessageTransformingHandler>.() -> Unit = {}) {
|
||||
fun transform(
|
||||
expression: String,
|
||||
endpointConfigurer: GenericEndpointSpec<MessageTransformingHandler>.() -> Unit = {}
|
||||
) {
|
||||
|
||||
this.delegate.transform(expression, endpointConfigurer)
|
||||
}
|
||||
@@ -323,8 +336,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* Populate the [MessageTransformingHandler] for the [MethodInvokingTransformer]
|
||||
* to invoke the service method at runtime.
|
||||
*/
|
||||
fun transform(service: Any, methodName: String?,
|
||||
endpointConfigurer: GenericEndpointSpec<MessageTransformingHandler>.() -> Unit) {
|
||||
fun transform(
|
||||
service: Any, methodName: String?,
|
||||
endpointConfigurer: GenericEndpointSpec<MessageTransformingHandler>.() -> Unit
|
||||
) {
|
||||
|
||||
this.delegate.transform(service, methodName, endpointConfigurer)
|
||||
}
|
||||
@@ -334,8 +349,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* [org.springframework.integration.handler.MessageProcessor] from provided [MessageProcessorSpec].
|
||||
* In addition, accept options for the integration endpoint using [GenericEndpointSpec].
|
||||
*/
|
||||
fun transform(messageProcessorSpec: MessageProcessorSpec<*>,
|
||||
endpointConfigurer: GenericEndpointSpec<MessageTransformingHandler>.() -> Unit = {}) {
|
||||
fun transform(
|
||||
messageProcessorSpec: MessageProcessorSpec<*>,
|
||||
endpointConfigurer: GenericEndpointSpec<MessageTransformingHandler>.() -> Unit = {}
|
||||
) {
|
||||
|
||||
this.delegate.transform(messageProcessorSpec, endpointConfigurer)
|
||||
}
|
||||
@@ -370,8 +387,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* the provided [MessageProcessorSpec].
|
||||
* In addition, accept options for the integration endpoint using [KotlinFilterEndpointSpec].
|
||||
*/
|
||||
fun filter(messageProcessorSpec: MessageProcessorSpec<*>,
|
||||
filterConfigurer: KotlinFilterEndpointSpec.() -> Unit = {}) {
|
||||
fun filter(
|
||||
messageProcessorSpec: MessageProcessorSpec<*>,
|
||||
filterConfigurer: KotlinFilterEndpointSpec.() -> Unit = {}
|
||||
) {
|
||||
|
||||
this.delegate.filter(messageProcessorSpec) { filterConfigurer(KotlinFilterEndpointSpec(it)) }
|
||||
}
|
||||
@@ -382,8 +401,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* In addition, accept options for the integration endpoint using [KotlinFilterEndpointSpec].
|
||||
* @since 5.3.1
|
||||
*/
|
||||
fun filter(messageSelector: MessageSelector,
|
||||
filterConfigurer: KotlinFilterEndpointSpec.() -> Unit = {}) {
|
||||
fun filter(
|
||||
messageSelector: MessageSelector,
|
||||
filterConfigurer: KotlinFilterEndpointSpec.() -> Unit = {}
|
||||
) {
|
||||
|
||||
this.delegate.filter(Message::class.java, messageSelector) { filterConfigurer(KotlinFilterEndpointSpec(it)) }
|
||||
}
|
||||
@@ -419,8 +440,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* to invoke the `method` for provided `bean` at runtime.
|
||||
* In addition, accept options for the integration endpoint using [GenericEndpointSpec].
|
||||
*/
|
||||
fun handle(beanName: String, methodName: String?,
|
||||
endpointConfigurer: GenericEndpointSpec<ServiceActivatingHandler>.() -> Unit) {
|
||||
fun handle(
|
||||
beanName: String, methodName: String?,
|
||||
endpointConfigurer: GenericEndpointSpec<ServiceActivatingHandler>.() -> Unit
|
||||
) {
|
||||
|
||||
this.delegate.handle(beanName, methodName, endpointConfigurer)
|
||||
}
|
||||
@@ -441,8 +464,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* to invoke the `method` for provided `bean` at runtime.
|
||||
* In addition, accept options for the integration endpoint using [GenericEndpointSpec].
|
||||
*/
|
||||
fun handle(service: Any, methodName: String?,
|
||||
endpointConfigurer: GenericEndpointSpec<ServiceActivatingHandler>.() -> Unit) {
|
||||
fun handle(
|
||||
service: Any, methodName: String?,
|
||||
endpointConfigurer: GenericEndpointSpec<ServiceActivatingHandler>.() -> Unit
|
||||
) {
|
||||
|
||||
this.delegate.handle(service, methodName, endpointConfigurer)
|
||||
}
|
||||
@@ -463,8 +488,9 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* In addition, accept options for the integration endpoint using [GenericEndpointSpec].
|
||||
*/
|
||||
inline fun <reified P> handle(
|
||||
crossinline handler: (P, MessageHeaders) -> Any,
|
||||
crossinline endpointConfigurer: GenericEndpointSpec<ServiceActivatingHandler>.() -> Unit) {
|
||||
crossinline handler: (P, MessageHeaders) -> Any,
|
||||
crossinline endpointConfigurer: GenericEndpointSpec<ServiceActivatingHandler>.() -> Unit
|
||||
) {
|
||||
|
||||
this.delegate.handle(P::class.java, { p, h -> handler(p, h) }) { endpointConfigurer(it) }
|
||||
}
|
||||
@@ -473,8 +499,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* Populate a [ServiceActivatingHandler] for the [MessageProcessor] from the provided [MessageProcessorSpec].
|
||||
* In addition, accept options for the integration endpoint using [GenericEndpointSpec].
|
||||
*/
|
||||
fun handle(messageProcessorSpec: MessageProcessorSpec<*>,
|
||||
endpointConfigurer: GenericEndpointSpec<ServiceActivatingHandler>.() -> Unit = {}) {
|
||||
fun handle(
|
||||
messageProcessorSpec: MessageProcessorSpec<*>,
|
||||
endpointConfigurer: GenericEndpointSpec<ServiceActivatingHandler>.() -> Unit = {}
|
||||
) {
|
||||
|
||||
this.delegate.handle(messageProcessorSpec, endpointConfigurer)
|
||||
}
|
||||
@@ -484,8 +512,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* [MessageHandler] implementation from `Namespace Factory`:
|
||||
* In addition, accept options for the integration endpoint using [GenericEndpointSpec].
|
||||
*/
|
||||
fun <H : MessageHandler> handle(messageHandlerSpec: MessageHandlerSpec<*, H>,
|
||||
endpointConfigurer: GenericEndpointSpec<H>.() -> Unit = {}) {
|
||||
fun <H : MessageHandler> handle(
|
||||
messageHandlerSpec: MessageHandlerSpec<*, H>,
|
||||
endpointConfigurer: GenericEndpointSpec<H>.() -> Unit = {}
|
||||
) {
|
||||
|
||||
this.delegate.handle(messageHandlerSpec, endpointConfigurer)
|
||||
}
|
||||
@@ -503,8 +533,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* [MessageHandler] lambda.
|
||||
* In addition, accept options for the integration endpoint using [GenericEndpointSpec].
|
||||
*/
|
||||
fun handle(messageHandler: (Message<*>) -> Unit,
|
||||
endpointConfigurer: GenericEndpointSpec<MessageHandler>.() -> Unit) {
|
||||
fun handle(
|
||||
messageHandler: (Message<*>) -> Unit,
|
||||
endpointConfigurer: GenericEndpointSpec<MessageHandler>.() -> Unit
|
||||
) {
|
||||
|
||||
this.delegate.handle(MessageHandler { messageHandler(it) }, endpointConfigurer)
|
||||
}
|
||||
@@ -547,8 +579,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* using header values from provided [MapBuilder].
|
||||
* In addition, accept options for the integration endpoint using [GenericEndpointSpec].
|
||||
*/
|
||||
fun enrichHeaders(headers: MapBuilder<*, String, Any>,
|
||||
endpointConfigurer: GenericEndpointSpec<MessageTransformingHandler>.() -> Unit = {}) {
|
||||
fun enrichHeaders(
|
||||
headers: MapBuilder<*, String, Any>,
|
||||
endpointConfigurer: GenericEndpointSpec<MessageTransformingHandler>.() -> Unit = {}
|
||||
) {
|
||||
|
||||
this.delegate.enrichHeaders(headers, endpointConfigurer)
|
||||
}
|
||||
@@ -559,8 +593,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* `values` can apply an [Expression]
|
||||
* to be evaluated against a request [Message].
|
||||
*/
|
||||
fun enrichHeaders(headers: Map<String, Any>,
|
||||
endpointConfigurer: GenericEndpointSpec<MessageTransformingHandler>.() -> Unit = {}) {
|
||||
fun enrichHeaders(
|
||||
headers: Map<String, Any>,
|
||||
endpointConfigurer: GenericEndpointSpec<MessageTransformingHandler>.() -> Unit = {}
|
||||
) {
|
||||
|
||||
this.delegate.enrichHeaders(headers, endpointConfigurer)
|
||||
}
|
||||
@@ -586,8 +622,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* Populate the [ExpressionEvaluatingSplitter] with provided
|
||||
* SpEL expression.
|
||||
*/
|
||||
fun split(expression: String,
|
||||
endpointConfigurer: KotlinSplitterEndpointSpec<ExpressionEvaluatingSplitter>.() -> Unit = {}) {
|
||||
fun split(
|
||||
expression: String,
|
||||
endpointConfigurer: KotlinSplitterEndpointSpec<ExpressionEvaluatingSplitter>.() -> Unit = {}
|
||||
) {
|
||||
|
||||
this.delegate.split(expression) { endpointConfigurer(KotlinSplitterEndpointSpec(it)) }
|
||||
}
|
||||
@@ -605,8 +643,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* `method` of the `bean` at runtime.
|
||||
* In addition, accept options for the integration endpoint using [KotlinSplitterEndpointSpec].
|
||||
*/
|
||||
fun split(service: Any, methodName: String?,
|
||||
splitterConfigurer: KotlinSplitterEndpointSpec<MethodInvokingSplitter>.() -> Unit) {
|
||||
fun split(
|
||||
service: Any, methodName: String?,
|
||||
splitterConfigurer: KotlinSplitterEndpointSpec<MethodInvokingSplitter>.() -> Unit
|
||||
) {
|
||||
|
||||
this.delegate.split(service, methodName) { splitterConfigurer(KotlinSplitterEndpointSpec(it)) }
|
||||
}
|
||||
@@ -624,8 +664,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* `method` of the `bean` at runtime.
|
||||
* In addition, accept options for the integration endpoint using [KotlinSplitterEndpointSpec].
|
||||
*/
|
||||
fun split(beanName: String, methodName: String?,
|
||||
splitterConfigurer: KotlinSplitterEndpointSpec<MethodInvokingSplitter>.() -> Unit) {
|
||||
fun split(
|
||||
beanName: String, methodName: String?,
|
||||
splitterConfigurer: KotlinSplitterEndpointSpec<MethodInvokingSplitter>.() -> Unit
|
||||
) {
|
||||
|
||||
this.delegate.split(beanName, methodName) { splitterConfigurer(KotlinSplitterEndpointSpec(it)) }
|
||||
}
|
||||
@@ -635,8 +677,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* [MessageProcessor] at runtime from provided [MessageProcessorSpec].
|
||||
* In addition, accept options for the integration endpoint using [KotlinSplitterEndpointSpec].
|
||||
*/
|
||||
fun split(messageProcessorSpec: MessageProcessorSpec<*>,
|
||||
splitterConfigurer: KotlinSplitterEndpointSpec<MethodInvokingSplitter>.() -> Unit = {}) {
|
||||
fun split(
|
||||
messageProcessorSpec: MessageProcessorSpec<*>,
|
||||
splitterConfigurer: KotlinSplitterEndpointSpec<MethodInvokingSplitter>.() -> Unit = {}
|
||||
) {
|
||||
|
||||
this.delegate.split(messageProcessorSpec) { splitterConfigurer(KotlinSplitterEndpointSpec(it)) }
|
||||
}
|
||||
@@ -644,8 +688,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
/**
|
||||
* Populate the provided [AbstractMessageSplitter] to the current integration flow position.
|
||||
*/
|
||||
fun <S : AbstractMessageSplitter> split(splitterMessageHandlerSpec: MessageHandlerSpec<*, S>,
|
||||
splitterConfigurer: KotlinSplitterEndpointSpec<S>.() -> Unit = {}) {
|
||||
fun <S : AbstractMessageSplitter> split(
|
||||
splitterMessageHandlerSpec: MessageHandlerSpec<*, S>,
|
||||
splitterConfigurer: KotlinSplitterEndpointSpec<S>.() -> Unit = {}
|
||||
) {
|
||||
|
||||
this.delegate.split(splitterMessageHandlerSpec) { splitterConfigurer(KotlinSplitterEndpointSpec(it)) }
|
||||
}
|
||||
@@ -654,8 +700,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* Populate the provided [AbstractMessageSplitter] to the current integration
|
||||
* flow position.
|
||||
*/
|
||||
fun <S : AbstractMessageSplitter> split(splitter: S,
|
||||
splitterConfigurer: KotlinSplitterEndpointSpec<S>.() -> Unit = {}) {
|
||||
fun <S : AbstractMessageSplitter> split(
|
||||
splitter: S,
|
||||
splitterConfigurer: KotlinSplitterEndpointSpec<S>.() -> Unit = {}
|
||||
) {
|
||||
|
||||
this.delegate.split(splitter) { splitterConfigurer(KotlinSplitterEndpointSpec(it)) }
|
||||
}
|
||||
@@ -671,8 +719,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* Populate the provided [MessageTransformingHandler] for the provided
|
||||
* [HeaderFilter].
|
||||
*/
|
||||
fun headerFilter(headerFilter: HeaderFilter,
|
||||
endpointConfigurer: GenericEndpointSpec<MessageTransformingHandler>.() -> Unit) {
|
||||
fun headerFilter(
|
||||
headerFilter: HeaderFilter,
|
||||
endpointConfigurer: GenericEndpointSpec<MessageTransformingHandler>.() -> Unit
|
||||
) {
|
||||
|
||||
this.delegate.headerFilter(headerFilter, endpointConfigurer)
|
||||
}
|
||||
@@ -682,8 +732,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* with provided [MessageStore].
|
||||
* In addition, accept options for the integration endpoint using [GenericEndpointSpec].
|
||||
*/
|
||||
fun claimCheckIn(messageStore: MessageStore,
|
||||
endpointConfigurer: GenericEndpointSpec<MessageTransformingHandler>.() -> Unit = {}) {
|
||||
fun claimCheckIn(
|
||||
messageStore: MessageStore,
|
||||
endpointConfigurer: GenericEndpointSpec<MessageTransformingHandler>.() -> Unit = {}
|
||||
) {
|
||||
|
||||
this.delegate.claimCheckIn(messageStore, endpointConfigurer)
|
||||
}
|
||||
@@ -701,8 +753,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* with provided [MessageStore] and `removeMessage` flag.
|
||||
* In addition, accept options for the integration endpoint using [GenericEndpointSpec].
|
||||
*/
|
||||
fun claimCheckOut(messageStore: MessageStore, removeMessage: Boolean,
|
||||
endpointConfigurer: GenericEndpointSpec<MessageTransformingHandler>.() -> Unit) {
|
||||
fun claimCheckOut(
|
||||
messageStore: MessageStore, removeMessage: Boolean,
|
||||
endpointConfigurer: GenericEndpointSpec<MessageTransformingHandler>.() -> Unit
|
||||
) {
|
||||
|
||||
this.delegate.claimCheckOut(messageStore, removeMessage, endpointConfigurer)
|
||||
}
|
||||
@@ -745,8 +799,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* Populate the [MethodInvokingRouter] for provided bean and its method
|
||||
* with provided options from [KotlinRouterSpec].
|
||||
*/
|
||||
fun route(beanName: String, method: String?,
|
||||
routerConfigurer: KotlinRouterSpec<Any, MethodInvokingRouter>.() -> Unit) {
|
||||
fun route(
|
||||
beanName: String, method: String?,
|
||||
routerConfigurer: KotlinRouterSpec<Any, MethodInvokingRouter>.() -> Unit
|
||||
) {
|
||||
|
||||
this.delegate.route(beanName, method) { routerConfigurer(KotlinRouterSpec(it)) }
|
||||
}
|
||||
@@ -763,8 +819,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* Populate the [MethodInvokingRouter] for the method
|
||||
* of the provided service and its method with provided options from [KotlinRouterSpec].
|
||||
*/
|
||||
fun route(service: Any, methodName: String?,
|
||||
routerConfigurer: KotlinRouterSpec<Any, MethodInvokingRouter>.() -> Unit) {
|
||||
fun route(
|
||||
service: Any, methodName: String?,
|
||||
routerConfigurer: KotlinRouterSpec<Any, MethodInvokingRouter>.() -> Unit
|
||||
) {
|
||||
|
||||
this.delegate.route(service, methodName) { routerConfigurer(KotlinRouterSpec(it)) }
|
||||
}
|
||||
@@ -773,8 +831,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* Populate the [ExpressionEvaluatingRouter] for provided SpEL expression
|
||||
* with provided options from [KotlinRouterSpec].
|
||||
*/
|
||||
fun <T> route(expression: String,
|
||||
routerConfigurer: KotlinRouterSpec<T, ExpressionEvaluatingRouter>.() -> Unit = {}) {
|
||||
fun <T> route(
|
||||
expression: String,
|
||||
routerConfigurer: KotlinRouterSpec<T, ExpressionEvaluatingRouter>.() -> Unit = {}
|
||||
) {
|
||||
|
||||
this.delegate.route<T>(expression) { routerConfigurer(KotlinRouterSpec(it)) }
|
||||
}
|
||||
@@ -783,8 +843,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* Populate the [MethodInvokingRouter] for the [MessageProcessor]
|
||||
* from the provided [MessageProcessorSpec] with default options.
|
||||
*/
|
||||
fun route(messageProcessorSpec: MessageProcessorSpec<*>,
|
||||
routerConfigurer: KotlinRouterSpec<Any, MethodInvokingRouter>.() -> Unit = {}) {
|
||||
fun route(
|
||||
messageProcessorSpec: MessageProcessorSpec<*>,
|
||||
routerConfigurer: KotlinRouterSpec<Any, MethodInvokingRouter>.() -> Unit = {}
|
||||
) {
|
||||
|
||||
this.delegate.route(messageProcessorSpec) { routerConfigurer(KotlinRouterSpec(it)) }
|
||||
}
|
||||
@@ -800,7 +862,8 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* Populate the [ErrorMessageExceptionTypeRouter] with options from the [KotlinRouterSpec].
|
||||
*/
|
||||
fun routeByException(
|
||||
routerConfigurer: KotlinRouterSpec<Class<out Throwable>, ErrorMessageExceptionTypeRouter>.() -> Unit) {
|
||||
routerConfigurer: KotlinRouterSpec<Class<out Throwable>, ErrorMessageExceptionTypeRouter>.() -> Unit
|
||||
) {
|
||||
|
||||
this.delegate.routeByException { routerConfigurer(KotlinRouterSpec(it)) }
|
||||
}
|
||||
@@ -852,12 +915,15 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* [org.springframework.integration.gateway.GatewayMessageHandler] for the
|
||||
* provided `subflow` with options from [GatewayEndpointSpec].
|
||||
*/
|
||||
fun gateway(endpointConfigurer: GatewayEndpointSpec.() -> Unit,
|
||||
flow: KotlinIntegrationFlowDefinition.() -> Unit) {
|
||||
fun gateway(
|
||||
endpointConfigurer: GatewayEndpointSpec.() -> Unit,
|
||||
flow: KotlinIntegrationFlowDefinition.() -> Unit
|
||||
) {
|
||||
|
||||
this.delegate.gateway(
|
||||
IntegrationFlow { flow(KotlinIntegrationFlowDefinition(it)) },
|
||||
Consumer(endpointConfigurer))
|
||||
IntegrationFlow { flow(KotlinIntegrationFlowDefinition(it)) },
|
||||
Consumer(endpointConfigurer)
|
||||
)
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1000,8 +1066,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* based on the provided [MessageChannel] for scattering function
|
||||
* and [AggregatorSpec] for gathering function.
|
||||
*/
|
||||
fun scatterGather(scatterChannel: MessageChannel, gatherer: AggregatorSpec.() -> Unit,
|
||||
scatterGather: ScatterGatherSpec.() -> Unit) {
|
||||
fun scatterGather(
|
||||
scatterChannel: MessageChannel, gatherer: AggregatorSpec.() -> Unit,
|
||||
scatterGather: ScatterGatherSpec.() -> Unit
|
||||
) {
|
||||
|
||||
this.delegate.scatterGather(scatterChannel, Consumer(gatherer), Consumer(scatterGather))
|
||||
}
|
||||
@@ -1022,7 +1090,7 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
*/
|
||||
fun scatterGather(scatterer: KotlinRecipientListRouterSpec.() -> Unit, gatherer: AggregatorSpec.() -> Unit) {
|
||||
this.delegate.scatterGather(Consumer { scatterer(KotlinRecipientListRouterSpec(it)) },
|
||||
Consumer { gatherer(it) })
|
||||
Consumer { gatherer(it) })
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1030,11 +1098,13 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* based on the provided [KotlinRecipientListRouterSpec] for scattering function
|
||||
* and [AggregatorSpec] for gathering function.
|
||||
*/
|
||||
fun scatterGather(scatterer: KotlinRecipientListRouterSpec.() -> Unit, gatherer: AggregatorSpec.() -> Unit,
|
||||
scatterGather: ScatterGatherSpec.() -> Unit) {
|
||||
fun scatterGather(
|
||||
scatterer: KotlinRecipientListRouterSpec.() -> Unit, gatherer: AggregatorSpec.() -> Unit,
|
||||
scatterGather: ScatterGatherSpec.() -> Unit
|
||||
) {
|
||||
|
||||
this.delegate.scatterGather(Consumer { scatterer(KotlinRecipientListRouterSpec(it)) },
|
||||
Consumer { gatherer(it) }, Consumer { scatterGather(it) })
|
||||
Consumer { gatherer(it) }, Consumer { scatterGather(it) })
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1050,8 +1120,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* Populate a [ServiceActivatingHandler] instance to perform [MessageTriggerAction]
|
||||
* and endpoint options from [GenericEndpointSpec].
|
||||
*/
|
||||
fun trigger(triggerActionId: String,
|
||||
endpointConfigurer: GenericEndpointSpec<ServiceActivatingHandler>.() -> Unit = {}) {
|
||||
fun trigger(
|
||||
triggerActionId: String,
|
||||
endpointConfigurer: GenericEndpointSpec<ServiceActivatingHandler>.() -> Unit = {}
|
||||
) {
|
||||
|
||||
this.delegate.trigger(triggerActionId, endpointConfigurer)
|
||||
}
|
||||
@@ -1060,8 +1132,10 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
* Populate a [ServiceActivatingHandler] instance to perform [MessageTriggerAction]
|
||||
* and endpoint options from [GenericEndpointSpec].
|
||||
*/
|
||||
fun trigger(triggerAction: MessageTriggerAction,
|
||||
endpointConfigurer: GenericEndpointSpec<ServiceActivatingHandler>.() -> Unit = {}) {
|
||||
fun trigger(
|
||||
triggerAction: MessageTriggerAction,
|
||||
endpointConfigurer: GenericEndpointSpec<ServiceActivatingHandler>.() -> Unit = {}
|
||||
) {
|
||||
|
||||
this.delegate.trigger(triggerAction, Consumer(endpointConfigurer))
|
||||
}
|
||||
@@ -1075,4 +1149,14 @@ class KotlinIntegrationFlowDefinition(@PublishedApi internal val delegate: Integ
|
||||
this.delegate.fluxTransform(fluxFunction)
|
||||
}
|
||||
|
||||
/**
|
||||
* Add one or more [ChannelInterceptor] implementations
|
||||
* to the current [MessageChannel], in the given order, after any interceptors already registered.
|
||||
* @param interceptorArray one or more [ChannelInterceptor]s.
|
||||
* @since 5.5.8
|
||||
*/
|
||||
fun intercept(vararg interceptorArray: ChannelInterceptor) {
|
||||
this.delegate.intercept(*interceptorArray)
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user