Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@
import io.serverlessworkflow.impl.WorkflowDefinition;
import io.serverlessworkflow.impl.WorkflowMutablePosition;
import io.serverlessworkflow.impl.executors.CallFunctionExecutorBuilder;
import io.serverlessworkflow.impl.executors.CallableTaskFactory;
import io.serverlessworkflow.impl.executors.CallableTask;
import java.util.Map;
import java.util.Optional;
import java.util.function.Consumer;
Expand All @@ -42,7 +42,7 @@ public int priority() {
}

@Override
public CallableTaskFactory init(
public CallableTask build(
CallFunction task, WorkflowDefinition definition, WorkflowMutablePosition position) {
if (CallJava.JAVA_CALL_KEY.equals(task.getCall())) {
if (task.getWith() == null) {
Expand All @@ -62,32 +62,30 @@ public CallableTaskFactory init(
Optional<Class<?>> output =
(Optional<Class<?>>) props.getOrDefault(CallJava.OUTPUT_CLASS_KEY, Optional.empty());
if (obj instanceof ContextFunction fn) {
return () -> new JavaContextFunctionCallExecutor(input, output, fn);
return new JavaContextFunctionCallExecutor(input, output, fn);
} else if (obj instanceof FilterFunction fn) {
return () -> new JavaFilterFunctionCallExecutor(input, output, fn);
return new JavaFilterFunctionCallExecutor(input, output, fn);
} else if (obj instanceof LoopFunction loop) {
return () ->
new JavaLoopFunctionCallExecutor(
loop, (String) props.get(CallJava.VAR_NAME_KEY), input, output);
return new JavaLoopFunctionCallExecutor(
loop, (String) props.get(CallJava.VAR_NAME_KEY), input, output);
} else if (obj instanceof LoopFunctionIndex loop) {
return () ->
new JavaLoopFunctionIndexCallExecutor(
loop,
(String) props.get(CallJava.VAR_NAME_KEY),
(String) props.get(CallJava.INDEX_NAME_KEY),
input,
output);
return new JavaLoopFunctionIndexCallExecutor(
loop,
(String) props.get(CallJava.VAR_NAME_KEY),
(String) props.get(CallJava.INDEX_NAME_KEY),
input,
output);

} else if (obj instanceof Function fn) {
return () -> new JavaFunctionCallExecutor(input, output, fn);
return new JavaFunctionCallExecutor(input, output, fn);
} else if (obj instanceof Consumer consumer) {
return () -> new JavaConsumerCallExecutor(input, consumer);
return new JavaConsumerCallExecutor(input, consumer);
} else {
throw new UnsupportedOperationException("Unrecognized function " + obj);
}
} else {
logger.info("Calling regular function handler for task call {}", task.getCall());
return super.init(task, definition, position);
return super.build(task, definition, position);
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -24,8 +24,8 @@
import io.serverlessworkflow.impl.WorkflowMutablePosition;
import io.serverlessworkflow.impl.WorkflowUtils;
import io.serverlessworkflow.impl.WorkflowValueResolver;
import io.serverlessworkflow.impl.executors.CallableTask;
import io.serverlessworkflow.impl.executors.CallableTaskBuilder;
import io.serverlessworkflow.impl.executors.CallableTaskFactory;
import java.net.URI;
import java.util.Map;
import java.util.Optional;
Expand All @@ -38,7 +38,7 @@ public boolean accept(Class<? extends TaskBase> clazz) {
}

@Override
public CallableTaskFactory init(
public CallableTask build(
CallA2A task, WorkflowDefinition definition, WorkflowMutablePosition position) {
A2AArguments args = task.getWith();

Expand Down Expand Up @@ -88,6 +88,6 @@ public CallableTaskFactory init(
parameters.getString(),
a2aParameters != null ? a2aParameters.getAdditionalProperties() : null));
}
return () -> new A2AExecutor(uriSupplier, dispatcher, mapResolver);
return new A2AExecutor(uriSupplier, dispatcher, mapResolver);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@
public class CallFunctionExecutorBuilder implements CallableTaskBuilder<CallFunction> {

@Override
public CallableTaskFactory init(
public CallableTask build(
CallFunction task, WorkflowDefinition definition, WorkflowMutablePosition position) {
String functionName = task.getCall();
Use use = definition.workflow().getUse();
Expand Down Expand Up @@ -86,7 +86,7 @@ catalogEndpoint, pathFromFunctionName(functionName.substring(0, indexOf))),
? WorkflowUtils.buildMapResolver(
definition.application(), functionArgs.getAdditionalProperties())
: (w, t, m) -> Map.of();
return () -> this.build(executorBuilder, args);
return buildCallable(executorBuilder, args);
}

private String pathFromFunctionName(String functionName) {
Expand Down Expand Up @@ -120,7 +120,7 @@ public boolean accept(Class<? extends TaskBase> clazz) {
return clazz.equals(CallFunction.class);
}

private CallableTask build(
private CallableTask buildCallable(
TaskExecutorBuilder<? extends TaskBase> executorBuilder, WorkflowValueResolver<?> args) {
TaskExecutor<? extends TaskBase> executor = executorBuilder.build();
return (w, t, m) ->
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,6 @@ public class CallTaskExecutor<T extends TaskBase> extends RegularTaskExecutor<T>

public static class CallTaskExecutorBuilder<T extends TaskBase>
extends RegularTaskExecutorBuilder<T> {
private CallableTaskFactory callableFactory;
private List<CallableTaskProxyBuilder> callableProxyBuilders;
private CallableTask callable;

Expand All @@ -44,12 +43,11 @@ protected CallTaskExecutorBuilder(
definition.application().callableProxyBuilders().stream()
.filter(t -> t.accept(task))
.toList();
this.callableFactory = callableBuilder.init(task, definition, position);
this.callable = callableBuilder.build(task, definition, position);
}

@Override
public CallTaskExecutor<T> buildInstance() {
this.callable = callableFactory.get();
for (CallableTaskProxyBuilder callableBuilder : callableProxyBuilders) {
this.callable = callableBuilder.build(callable);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,5 +24,5 @@ public interface CallableTaskBuilder<T extends TaskBase> extends ServicePriority

boolean accept(Class<? extends TaskBase> clazz);

CallableTaskFactory init(T task, WorkflowDefinition definition, WorkflowMutablePosition position);
CallableTask build(T task, WorkflowDefinition definition, WorkflowMutablePosition position);
}

This file was deleted.

Original file line number Diff line number Diff line change
Expand Up @@ -24,8 +24,8 @@
import io.serverlessworkflow.impl.WorkflowDefinition;
import io.serverlessworkflow.impl.WorkflowMutablePosition;
import io.serverlessworkflow.impl.WorkflowUtils;
import io.serverlessworkflow.impl.executors.CallableTask;
import io.serverlessworkflow.impl.executors.CallableTaskBuilder;
import io.serverlessworkflow.impl.executors.CallableTaskFactory;
import java.util.Map;
import java.util.Objects;

Expand All @@ -37,7 +37,7 @@ public boolean accept(Class<? extends TaskBase> clazz) {
}

@Override
public CallableTaskFactory init(
public CallableTask build(
CallGRPC task, WorkflowDefinition definition, WorkflowMutablePosition position) {
GRPCArguments with = task.getWith();
WithGRPCService service = with.getService();
Expand All @@ -53,16 +53,13 @@ public CallableTaskFactory init(
Objects.requireNonNull(
serviceDescriptor.findMethodByName(with.getMethod()),
"Method not found: " + with.getMethod());
return () ->
new GrpcExecutor(
service.getHost(),
service.getPort(),
WorkflowUtils.buildMapResolver(
definition.application(),
with.getArguments() != null
? with.getArguments().getAdditionalProperties()
: Map.of()),
serviceDescriptor,
methodDescriptor);
return new GrpcExecutor(
service.getHost(),
service.getPort(),
WorkflowUtils.buildMapResolver(
definition.application(),
with.getArguments() != null ? with.getArguments().getAdditionalProperties() : Map.of()),
serviceDescriptor,
methodDescriptor);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -24,14 +24,14 @@
import io.serverlessworkflow.impl.WorkflowDefinition;
import io.serverlessworkflow.impl.WorkflowMutablePosition;
import io.serverlessworkflow.impl.WorkflowValueResolver;
import io.serverlessworkflow.impl.executors.CallableTask;
import io.serverlessworkflow.impl.executors.CallableTaskBuilder;
import io.serverlessworkflow.impl.executors.CallableTaskFactory;
import java.net.URI;

public class CallableTaskHttpExecutorBuilder implements CallableTaskBuilder<CallHTTP> {

@Override
public CallableTaskFactory init(
public CallableTask build(
CallHTTP task, WorkflowDefinition definition, WorkflowMutablePosition position) {

HttpExecutorBuilder builder = HttpExecutorBuilder.builder(definition);
Expand Down Expand Up @@ -67,7 +67,7 @@ public CallableTaskFactory init(
builder.withBody(httpArgs.getBody());
builder.withMethod(httpArgs.getMethod().toUpperCase());
builder.redirect(httpArgs.isRedirect());
return () -> builder.build(uriSupplier);
return builder.build(uriSupplier);
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,8 +21,8 @@
import io.serverlessworkflow.api.types.TaskBase;
import io.serverlessworkflow.impl.WorkflowDefinition;
import io.serverlessworkflow.impl.WorkflowMutablePosition;
import io.serverlessworkflow.impl.executors.CallableTask;
import io.serverlessworkflow.impl.executors.CallableTaskBuilder;
import io.serverlessworkflow.impl.executors.CallableTaskFactory;
import io.serverlessworkflow.impl.executors.http.HttpExecutorBuilder;
import java.util.Map;

Expand All @@ -34,7 +34,7 @@ public boolean accept(Class<? extends TaskBase> clazz) {
}

@Override
public CallableTaskFactory init(
public CallableTask build(
CallOpenAPI task, WorkflowDefinition definition, WorkflowMutablePosition position) {
OpenAPIArguments with = task.getWith();
OpenAPIProcessor processor = new OpenAPIProcessor(with.getOperationId());
Expand All @@ -47,6 +47,6 @@ public CallableTaskFactory init(
HttpExecutorBuilder.builder(definition)
.withAuth(with.getAuthentication())
.redirect(with.isRedirect());
return () -> new OpenAPIExecutor(processor, resource, parameters, builder);
return new OpenAPIExecutor(processor, resource, parameters, builder);
}
}
Loading