diff --git a/conductor-client/src/main/java/com/netflix/conductor/sdk/workflow/def/tasks/Event.java b/conductor-client/src/main/java/com/netflix/conductor/sdk/workflow/def/tasks/Event.java index 4911a635d..c5279e559 100644 --- a/conductor-client/src/main/java/com/netflix/conductor/sdk/workflow/def/tasks/Event.java +++ b/conductor-client/src/main/java/com/netflix/conductor/sdk/workflow/def/tasks/Event.java @@ -22,6 +22,8 @@ public class Event extends Task { private static final String SINK_PARAMETER = "sink"; + private boolean asyncComplete; + /** * @param taskReferenceName Unique reference name within the workflow * @param eventSink qualified name of the event sink where the message is published. Using the @@ -38,6 +40,21 @@ public Event(String taskReferenceName, String eventSink) { Event(WorkflowTask workflowTask) { super(workflowTask); + this.asyncComplete = Boolean.TRUE.equals(workflowTask.isAsyncComplete()); + } + + public Event asyncComplete(boolean asyncComplete) { + this.asyncComplete = asyncComplete; + return this; + } + + public boolean isAsyncComplete() { + return asyncComplete; + } + + @Override + protected void updateWorkflowTask(WorkflowTask workflowTask) { + workflowTask.setAsyncComplete(asyncComplete); } public String getSink() { diff --git a/conductor-client/src/main/java/com/netflix/conductor/sdk/workflow/def/tasks/Http.java b/conductor-client/src/main/java/com/netflix/conductor/sdk/workflow/def/tasks/Http.java index 37468ff1c..3cb133b37 100644 --- a/conductor-client/src/main/java/com/netflix/conductor/sdk/workflow/def/tasks/Http.java +++ b/conductor-client/src/main/java/com/netflix/conductor/sdk/workflow/def/tasks/Http.java @@ -34,6 +34,8 @@ public class Http extends Task { private final ObjectMapper objectMapper = new ObjectMapperProvider().getObjectMapper(); + private boolean asyncComplete; + private Input httpRequest; public Http(String taskReferenceName) { @@ -45,6 +47,7 @@ public Http(String taskReferenceName) { Http(WorkflowTask workflowTask) { super(workflowTask); + this.asyncComplete = Boolean.TRUE.equals(workflowTask.isAsyncComplete()); Object inputRequest = workflowTask.getInputParameters().get(INPUT_PARAM); if (inputRequest != null) { @@ -56,6 +59,15 @@ public Http(String taskReferenceName) { } } + public Http asyncComplete(boolean asyncComplete) { + this.asyncComplete = asyncComplete; + return this; + } + + public boolean isAsyncComplete() { + return asyncComplete; + } + public Http input(Input httpRequest) { this.httpRequest = httpRequest; return this; @@ -92,6 +104,7 @@ public Input getHttpRequest() { @Override protected void updateWorkflowTask(WorkflowTask workflowTask) { + workflowTask.setAsyncComplete(asyncComplete); workflowTask.getInputParameters().put(INPUT_PARAM, httpRequest); } diff --git a/conductor-client/src/test/java/com/netflix/conductor/sdk/workflow/def/TaskConversionsTests.java b/conductor-client/src/test/java/com/netflix/conductor/sdk/workflow/def/TaskConversionsTests.java index e0d044729..1ced27725 100644 --- a/conductor-client/src/test/java/com/netflix/conductor/sdk/workflow/def/TaskConversionsTests.java +++ b/conductor-client/src/test/java/com/netflix/conductor/sdk/workflow/def/TaskConversionsTests.java @@ -460,6 +460,32 @@ public void testHttpConverter() { System.out.println(taskFromWorkflowTask.getInput()); } + @Test + public void testHttpAsyncComplete() { + Http httpTask = new Http("http_ref"); + httpTask.asyncComplete(true); + + WorkflowTask workflowTask = httpTask.getWorkflowDefTasks().get(0); + assertTrue(workflowTask.isAsyncComplete()); + + Task fromWorkflowTask = TaskRegistry.getTask(workflowTask); + assertTrue(fromWorkflowTask instanceof Http); + assertTrue(((Http) fromWorkflowTask).isAsyncComplete()); + } + + @Test + public void testEventAsyncComplete() { + Event eventTask = new Event("event_ref", "sqs:my-queue"); + eventTask.asyncComplete(true); + + WorkflowTask workflowTask = eventTask.getWorkflowDefTasks().get(0); + assertTrue(workflowTask.isAsyncComplete()); + + Task fromWorkflowTask = TaskRegistry.getTask(workflowTask); + assertTrue(fromWorkflowTask instanceof Event); + assertTrue(((Event) fromWorkflowTask).isAsyncComplete()); + } + @Test public void testJQTaskConversion() { JQ jqTask = new JQ("task_name", "{ key3: (.key1.value1 + .key2.value2) }");