Skip to content
5 changes: 5 additions & 0 deletions changelogs/feature/process-step-orchestration.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
{
"author": "Nemikor",
"pullrequestId": 361,
"message": "Add a process step with orchestration context available in scenario and process orchestration"
}
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,13 @@ public static Boolean handleConditional(
return new ConditionalHandler(stepResultCloner, conditional).executeCondition(exchange);
}

public static void handleProcess(
final Exchange exchange,
final Optional<StepResultCloner> stepResultCloner,
final Optional<CompositeProcessTransformer> processTransformer) {
new ProcessHandler(stepResultCloner, processTransformer).processTransformation(exchange);
}

public static int handleIterations(
final Exchange exchange, final Optional<CompositeProcessStepIterations> iterations) {
return new IterationsHandler(iterations).determineIterations(exchange);
Expand Down Expand Up @@ -180,6 +187,20 @@ public boolean executeCondition(final Exchange exchange) {
}
}

@RequiredArgsConstructor(access = AccessLevel.PRIVATE)
static class ProcessHandler {
private final Optional<StepResultCloner> stepResultCloner;
private final Optional<CompositeProcessTransformer> processor;

@Handler
public void processTransformation(final Exchange exchange) {
final CompositeProcessOrchestrationContext context = retrieveOrchestrationContext(exchange);
if (processor.isPresent()) {
processor.get().process(context);
}
}
}

@RequiredArgsConstructor(access = AccessLevel.PRIVATE)
static class IterationsHandler {
private final Optional<CompositeProcessStepIterations> iterations;
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
package one.x1f.sip.foundation.core.declarative.orchestration.process;

/** Interface to expose orchestration context in process step */
@FunctionalInterface
public interface CompositeProcessTransformer {

/**
* Define processing on orchestration context
*
* @param context The current orchestration context
*/
void process(final CompositeProcessOrchestrationContext context);
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
package one.x1f.sip.foundation.core.declarative.orchestration.process.dsl;

import java.util.Optional;
import lombok.AccessLevel;
import lombok.Getter;
import one.x1f.sip.foundation.core.declarative.orchestration.process.CompositeProcessTransformer;
import one.x1f.sip.foundation.core.declarative.process.CompositeProcessDefinition;

/**
* DSL class used for construction conditional calls after main condition
*
* @param <R> DSL handle for the return DSL Verb/type.
*/
public final class CallProcess<R> extends ProcessDslBase<CallProcess<R>, R>
implements CallableWithinProcessDefinition {

@Getter(AccessLevel.PACKAGE)
private Optional<CompositeProcessTransformer> process = Optional.empty();

CallProcess(R dslReturnDefinition, CompositeProcessDefinition compositeProcess) {
super(dslReturnDefinition, compositeProcess);
}

R process(final CompositeProcessTransformer expression) {
process = Optional.of(expression);
return getDslReturnDefinition();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -6,9 +6,7 @@
import lombok.Getter;
import lombok.experimental.Delegate;
import one.x1f.sip.foundation.core.declarative.orchestration.common.dsl.EndOfDsl;
import one.x1f.sip.foundation.core.declarative.orchestration.process.CompositeProcessStepConditional;
import one.x1f.sip.foundation.core.declarative.orchestration.process.CompositeProcessStepIterations;
import one.x1f.sip.foundation.core.declarative.orchestration.process.CompositeProcessStepSplitExpression;
import one.x1f.sip.foundation.core.declarative.orchestration.process.*;
import one.x1f.sip.foundation.core.declarative.process.CompositeProcessDefinition;

/** DSL class for specifying orchestration of complex processes with or without conditions */
Expand Down Expand Up @@ -110,4 +108,11 @@ public ProcessOrchestrationDefinition(final CompositeProcessDefinition composite
steps.add(def);
return def.parallelSplit(expression);
}

public ProcessOrchestrationDefinition process(CompositeProcessTransformer requestPreparation) {
CallProcess<ProcessOrchestrationDefinition> def =
new CallProcess(self(), getCompositeProcess());
steps.add(def);
return def.process(requestPreparation);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
import one.x1f.sip.foundation.core.declarative.orchestration.common.dsl.StepResultCloner;
import one.x1f.sip.foundation.core.declarative.orchestration.process.CompositeProcessStepRequestExtractor;
import one.x1f.sip.foundation.core.declarative.orchestration.process.CompositeProcessStepResponseConsumer;
import one.x1f.sip.foundation.core.declarative.orchestration.process.CompositeProcessTransformer;
import one.x1f.sip.foundation.core.declarative.scenario.IntegrationScenarioDefinition;

/**
Expand Down Expand Up @@ -39,6 +40,10 @@ public static Optional<CompositeProcessStepResponseConsumer> getResponseConsumer
return element.getResponseConsumer();
}

public static Optional<CompositeProcessTransformer> getProcess(CallProcess element) {
return element.getProcess();
}

public static List<CallNestedCondition.ProcessBranchStatements> getConditionalStatements(
CallNestedCondition element) {
return element.getConditionalStatements();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,9 @@ void generateRoutes(final RoutesDefinition routesDefinition) {
new RouteGeneratorForSplitProcessConsumer(
getOrchestrationInfo(), ele, unhandledProcessConsumers)
.generateRoute(routeDef);
} else if (element instanceof CallProcess<?> ele) {
new RouteGeneratorForProcessTransformer(getOrchestrationInfo(), ele)
.generateRoute(routeDef);
}
}

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,46 @@
package one.x1f.sip.foundation.core.declarative.orchestration.process.routebuilding;

import java.util.Optional;
import lombok.extern.slf4j.Slf4j;
import one.x1f.sip.foundation.core.declarative.orchestration.process.CompositeProcessOrchestrationHandlers;
import one.x1f.sip.foundation.core.declarative.orchestration.process.CompositeProcessOrchestrationInfo;
import one.x1f.sip.foundation.core.declarative.orchestration.process.CompositeProcessTransformer;
import one.x1f.sip.foundation.core.declarative.orchestration.process.dsl.CallProcess;
import one.x1f.sip.foundation.core.declarative.orchestration.process.dsl.RouteGeneratorInternalHelper;
import one.x1f.sip.foundation.core.util.exception.SIPFrameworkInitializationException;
import org.apache.camel.model.ProcessorDefinition;

/**
* Class for generating Camel routes for process consumer calls from a DSL
*
* <p><em>For internal use only</em>
*/
@Slf4j
@SuppressWarnings("rawtypes")
final class RouteGeneratorForProcessTransformer extends RouteGeneratorProcessBase {

private final CallProcess<?> definitionElement;

RouteGeneratorForProcessTransformer(
final CompositeProcessOrchestrationInfo orchestrationInfo,
final CallProcess definitionElement) {
super(orchestrationInfo);
this.definitionElement = definitionElement;
}

<T extends ProcessorDefinition<T>> void generateRoute(final T routeDefinition) {
Optional<CompositeProcessTransformer> compositeProcessTransformer =
RouteGeneratorInternalHelper.getProcess(definitionElement);
if (compositeProcessTransformer.isEmpty()) {
throw SIPFrameworkInitializationException.init(
"Empty process statement attached in orchestration for composite process '%s'",
getCompositeProcessId());
}
routeDefinition.process(
exchange ->
CompositeProcessOrchestrationHandlers.handleProcess(
exchange,
Optional.empty(),
RouteGeneratorInternalHelper.getProcess(definitionElement)));
}
}
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package one.x1f.sip.foundation.core.declarative.orchestration.scenario;

import java.util.*;
import java.util.function.Consumer;
import java.util.function.Predicate;
import lombok.AccessLevel;
import lombok.NoArgsConstructor;
Expand Down Expand Up @@ -50,6 +51,11 @@ public static <M> ContextPredicateHandler<M> handleContextPredicate(
return new ContextPredicateHandler<>(predicate);
}

public static <M> void handleContextConsumer(
final Consumer<ScenarioOrchestrationContext<M>> consumer, Exchange exchange) {
new ContextConsumerHandler<>(consumer).acceptConsumer(exchange);
}

public static ThrowErrorOnUnhandledRequestHandler handleErrorThrownIfNoConsumerWasCalled() {
return new ThrowErrorOnUnhandledRequestHandler();
}
Expand Down Expand Up @@ -143,6 +149,16 @@ public boolean testPredicate(final Exchange exchange) {
}
}

@RequiredArgsConstructor(access = AccessLevel.PRIVATE)
static class ContextConsumerHandler<M> {
private final Consumer<ScenarioOrchestrationContext<M>> consumer;

public void acceptConsumer(final Exchange exchange) {
final ScenarioOrchestrationContext<M> context = retrieveOrchestrationContext(exchange);
consumer.accept(context);
}
}

@NoArgsConstructor(access = AccessLevel.PRIVATE)
static class ThrowErrorOnUnhandledRequestHandler {
@Handler
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,4 +8,5 @@ sealed interface CallableWithinProviderDefinition
permits CallScenarioConsumerBaseDefinition,
ConditionalCallScenarioConsumerDefinition,
ForLoopCallScenarioConsumerDefinition,
ProcessCallScenarioConsumerDefinition,
WhileLoopCallScenarioConsumerDefinition {}
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@

import java.util.ArrayList;
import java.util.List;
import java.util.function.Consumer;
import java.util.function.Predicate;
import lombok.AccessLevel;
import lombok.Getter;
Expand Down Expand Up @@ -124,4 +125,11 @@ public CallScenarioConsumerCatchAllDefinition<R, M> callAnyUnspecifiedScenarioCo
public R endDefinitionForThisProvider() {
return getDslReturnDefinition();
}

public S process(final Consumer<ScenarioOrchestrationContext<M>> consumer) {
ProcessCallScenarioConsumerDefinition<S, M> def =
new ProcessCallScenarioConsumerDefinition<>(self(), getIntegrationScenario());
nodes.add(def);
return def.process(consumer);
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
package one.x1f.sip.foundation.core.declarative.orchestration.scenario.dsl;

import java.util.function.Consumer;
import lombok.Getter;
import one.x1f.sip.foundation.core.declarative.orchestration.scenario.ScenarioOrchestrationContext;
import one.x1f.sip.foundation.core.declarative.scenario.IntegrationScenarioDefinition;

public final class ProcessCallScenarioConsumerDefinition<R, M>
extends ScenarioDslDefinitionBase<ProcessCallScenarioConsumerDefinition<R, M>, R, M>
implements CallableWithinProviderDefinition {
@Getter Consumer<ScenarioOrchestrationContext<M>> consumer;

ProcessCallScenarioConsumerDefinition(
final R dslReturnDefinition, final IntegrationScenarioDefinition integrationScenario) {
super(dslReturnDefinition, integrationScenario);
}

R process(Consumer<ScenarioOrchestrationContext<M>> consumer) {
this.consumer = consumer;
return getDslReturnDefinition();
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
package one.x1f.sip.foundation.core.declarative.orchestration.scenario.dsl;

import lombok.extern.slf4j.Slf4j;
import one.x1f.sip.foundation.core.declarative.orchestration.scenario.ScenarioOrchestrationHandlers;
import one.x1f.sip.foundation.core.declarative.orchestration.scenario.ScenarioOrchestrationInfo;
import org.apache.camel.model.ProcessorDefinition;

@SuppressWarnings("rawtypes")
@Slf4j
final class RouteGeneratorForProcessCallScenarioConsumerDefinition<M> extends RouteGeneratorBase {
private final ProcessCallScenarioConsumerDefinition<?, M> processDefinition;

RouteGeneratorForProcessCallScenarioConsumerDefinition(
final ScenarioOrchestrationInfo orchestrationInfo,
final ProcessCallScenarioConsumerDefinition<?, M> processDefinition) {
super(orchestrationInfo);
this.processDefinition = processDefinition;
}

<T extends ProcessorDefinition<T>> void generateRoute(final T routeDefinition) {
routeDefinition.process(
exchange ->
ScenarioOrchestrationHandlers.handleContextConsumer(
processDefinition.getConsumer(), exchange));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -142,6 +142,10 @@ void generateRoutes(final RoutesDefinition routesDefinition) {
new RouteGeneratorForForLoopCallScenarioConsumerDefinition<M>(
getOrchestrationInfo(), loopDef, overallUnhandledScenarioConsumers)
.generateRoute(routeDef);
} else if (element instanceof ProcessCallScenarioConsumerDefinition processDef) {
new RouteGeneratorForProcessCallScenarioConsumerDefinition<M>(
getOrchestrationInfo(), processDef)
.generateRoute(routeDef);
} else {
throw SIPFrameworkInitializationException.init(
"No handling defined for type %s used in orchestration for scenario %s",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -120,7 +120,14 @@ public class GetCustomerDebtByNameOrchestrator extends CompositeProcessBase {
public Orchestrator<CompositeProcessOrchestrationInfo> getOrchestrator() {
return ProcessOrchestrator.forOrchestrationDsl(
dsl -> {
dsl.callConsumer(getPartnerByName.class)
dsl.process(
context -> {
context
.getExchange()
.getMessage()
.setHeader("headerKey", "Header set in process");
})
.callConsumer(getPartnerByName.class)
.withNoResponseHandling()
.callConsumer(getPartnerDebtById.class)
.withRequestPreparation(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
import one.x1f.sip.foundation.core.declarative.annotation.InboundConnector;
import one.x1f.sip.foundation.core.declarative.annotation.IntegrationScenario;
import one.x1f.sip.foundation.core.declarative.annotation.OutboundConnector;
import one.x1f.sip.foundation.core.declarative.annotation.connector.extension.ResponseProcessor;
import one.x1f.sip.foundation.core.declarative.connector.GenericInboundConnectorBase;
import one.x1f.sip.foundation.core.declarative.connector.GenericOutboundConnectorBase;
import one.x1f.sip.foundation.core.declarative.orchestration.Orchestrator;
Expand All @@ -19,6 +20,7 @@
import one.x1f.sip.foundation.core.declarative.orchestration.scenario.ScenarioOrchestrationInfo;
import one.x1f.sip.foundation.core.declarative.orchestration.scenario.ScenarioOrchestrator;
import one.x1f.sip.foundation.core.declarative.scenario.IntegrationScenarioBase;
import org.apache.camel.Exchange;
import org.apache.camel.builder.EndpointConsumerBuilder;
import org.apache.camel.builder.EndpointProducerBuilder;
import org.apache.camel.builder.endpoint.StaticEndpointBuilders;
Expand Down Expand Up @@ -69,6 +71,13 @@ public Orchestrator<ScenarioOrchestrationInfo> getOrchestrator() {
dsl -> {
// first connector only calls first outbound connector
dsl.forInboundConnectors(InboundConnectorOne.class, InboundConnectorTwo.class)
.process(
context -> {
context
.getExchange()
.getMessage()
.setHeader("headerKey", "Header set in scenario orchestration");
})
.callOutboundConnector(OutboundConnectorOne.ID)
.withRequestPreparation(
context -> context.getOriginalRequest() + "-scenarioprepared")
Expand Down Expand Up @@ -142,6 +151,11 @@ public class InboundConnectorOne extends GenericInboundConnectorBase {
protected EndpointConsumerBuilder defineInitiatingEndpoint() {
return StaticEndpointBuilders.direct("dummyInputOne");
}

@ResponseProcessor
public void handleResponse(Exchange response) {
response.getMessage();
}
}

@InboundConnector(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -104,5 +104,6 @@ void WHEN_callingProcessOrchestratorInboundConnectors_THEN_ReceiveResponse() {
assertThat(response.getAmount()).isEqualTo(new BigDecimal("100000.00"));
assertThat(response.getRequestedPartnerName()).isEqualTo("MyOrchestratedPartner");
assertThat(response.getRequestedBy()).isEqualTo("Process Orchestrator");
assertThat(exchange.getMessage().getHeader("headerKey")).isEqualTo("Header set in process");
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,8 @@ void WHEN_callingFirstOrSecondInboundConnector_THEN_OnlyFirstOutboundConnectorRe
assertThat(responseFirstConnector).isInstanceOf(ScenarioResponse.class);
assertThat(responseFirstConnector.getValue()).isEqualTo(1);
assertThat(responseSecondConnector).isEqualTo(responseFirstConnector);
assertThat(exchangeFirstConnector.getMessage().getHeader("headerKey"))
.isEqualTo("Header set in scenario orchestration");
}

@Test
Expand Down
Loading