diff --git a/src/common/types.ts b/src/common/types.ts index 5941bb50..61792b14 100644 --- a/src/common/types.ts +++ b/src/common/types.ts @@ -74,6 +74,7 @@ export interface EventTaskDef extends CommonTaskDef { type: TaskType.EVENT; sink: string; asyncComplete?: boolean; + optional?: boolean; } export interface ForkJoinTaskDef extends CommonTaskDef { @@ -120,6 +121,7 @@ export interface HttpTaskDef extends CommonTaskDef { }; type: TaskType.HTTP; asyncComplete?: boolean; + optional?: boolean; } export interface InlineTaskInputParameters { @@ -131,6 +133,7 @@ export interface InlineTaskInputParameters { export interface InlineTaskDef extends CommonTaskDef { type: TaskType.INLINE; inputParameters: InlineTaskInputParameters; + optional?: boolean; } interface ContainingQueryExpression { @@ -141,6 +144,7 @@ interface ContainingQueryExpression { export interface JsonJQTransformTaskDef extends CommonTaskDef { type: TaskType.JSON_JQ_TRANSFORM; inputParameters: ContainingQueryExpression; + optional?: boolean; } export interface KafkaPublishInputParameters { @@ -157,16 +161,19 @@ export interface KafkaPublishTaskDef extends CommonTaskDef { kafka_request: KafkaPublishInputParameters; }; type: TaskType.KAFKA_PUBLISH; + optional?: boolean; } export interface SetVariableTaskDef extends CommonTaskDef { type: TaskType.SET_VARIABLE; inputParameters: Record; + optional?: boolean; } export interface SimpleTaskDef extends CommonTaskDef { type: TaskType.SIMPLE; inputParameters?: Record; + optional?: boolean; } export interface SubWorkflowTaskDef extends CommonTaskDef { @@ -177,6 +184,7 @@ export interface SubWorkflowTaskDef extends CommonTaskDef { version?: number; taskToDomain?: Record; }; + optional?: boolean; } export interface SwitchTaskDef extends CommonTaskDef { @@ -186,6 +194,7 @@ export interface SwitchTaskDef extends CommonTaskDef { defaultCase: TaskDefTypes[]; evaluatorType: "value-param" | "javascript"; expression: string; + optional?: boolean; } export interface TerminateTaskDef extends CommonTaskDef { @@ -196,7 +205,6 @@ export interface TerminateTaskDef extends CommonTaskDef { }; type: TaskType.TERMINATE; startDelay?: number; - optional?: boolean; } export interface WaitTaskDef extends CommonTaskDef { @@ -205,6 +213,7 @@ export interface WaitTaskDef extends CommonTaskDef { duration?: string; until?: string; }; + optional?: boolean; } export interface WorkflowDef diff --git a/src/core/__test__/executor.test.ts b/src/core/__test__/executor.test.ts index 44ac48dc..d1bc6d5b 100644 --- a/src/core/__test__/executor.test.ts +++ b/src/core/__test__/executor.test.ts @@ -145,6 +145,29 @@ describe("Executor", () => { expect(workflowStatusAfter.tasks?.[0]?.status).toEqual("COMPLETED"); }); + + test("Should run workflow with an optional http task", async () => { + const executor = new WorkflowExecutor(await clientPromise); + + await executor.registerWorkflow(true, { + name: "test_jssdk_workflow_with_optional_http_task", + version: 1, + ownerEmail: "developers@orkes.io", + tasks: [httpTask("test_jssdk_optional_http_task", { uri: "uncorrect_uri", method: "GET" }, false, true)], + inputParameters: [], + outputParameters: {}, + timeoutSeconds: 300, + }); + + const executionId = await executor.startWorkflow({ + name: "test_jssdk_workflow_with_optional_http_task", + input: {}, + version: 1, + }); + + const workflowStatus = await TestUtil.waitForWorkflowStatus(executor, executionId, "COMPLETED"); + expect(["FAILED", "COMPLETED_WITH_ERRORS"]).toContain(workflowStatus.tasks?.[0]?.status); + }); }); describe("Execute with Return Strategy and Consistency", () => { diff --git a/src/core/generators/ForkJoin.ts b/src/core/generators/ForkJoin.ts index 7ab9110d..05ae5f70 100644 --- a/src/core/generators/ForkJoin.ts +++ b/src/core/generators/ForkJoin.ts @@ -19,8 +19,6 @@ export const generateJoinTask = ( ...nameTaskNameGenerator("join", overrides), inputParameters: {}, joinOn: [], - optional: false, - asyncComplete: false, ...overrides, type: TaskType.JOIN, }); diff --git a/src/core/generators/TerminateTask.ts b/src/core/generators/TerminateTask.ts index f77b9176..7c585c96 100644 --- a/src/core/generators/TerminateTask.ts +++ b/src/core/generators/TerminateTask.ts @@ -17,7 +17,6 @@ export const generateTerminateTask = ( workflowOutput: {}, }, startDelay: 0, - optional: false, ...overrides, type: TaskType.TERMINATE, }); diff --git a/src/core/sdk/__test__/factory.test.ts b/src/core/sdk/__test__/factory.test.ts index 175cf027..3583ad19 100644 --- a/src/core/sdk/__test__/factory.test.ts +++ b/src/core/sdk/__test__/factory.test.ts @@ -114,7 +114,7 @@ describe("forkTask", () => { const tname = "forkTaskJoin"; const [forkTask, joinTask] = forkTaskJoin(tname, [ eventTask(tname, "prefix", "suffix"), - ]); + ], true); expect(forkTask).toEqual({ taskReferenceName: "forkTaskJoin", name: "forkTaskJoin", @@ -135,8 +135,7 @@ describe("forkTask", () => { taskReferenceName: "forkTaskJoin_join_ref", inputParameters: {}, joinOn: [], - optional: false, - asyncComplete: false, + optional: true, type: "JOIN", }); }); @@ -150,7 +149,6 @@ describe("httpTask", () => { method: "GET", }); expect(httpTaskObj).toEqual({ - asyncComplete: false, name: "httpTask", taskReferenceName: "httpTask", inputParameters: { @@ -159,6 +157,33 @@ describe("httpTask", () => { type: "HTTP", }); }); + + it("Should create an http task with asyncComplete property", () => { + const httpTaskObj = httpTask("testHttp", { uri: "https://example.com", method: "GET" }, true); + expect(httpTaskObj).toEqual({ + name: "testHttp", + taskReferenceName: "testHttp", + inputParameters: { + http_request: { uri: "https://example.com", method: "GET" }, + }, + asyncComplete: true, + type: "HTTP", + }); + }); + + it("Should create an http task with optional property", () => { + const httpTaskObj = httpTask("testHttp", { uri: "https://example.com", method: "GET" }, false, true); + expect(httpTaskObj).toEqual({ + name: "testHttp", + taskReferenceName: "testHttp", + inputParameters: { + http_request: { uri: "https://example.com", method: "GET" }, + }, + asyncComplete: false, + type: "HTTP", + optional: true, + }); + }); }); describe("inlineTask", () => { diff --git a/src/core/sdk/doWhile.ts b/src/core/sdk/doWhile.ts index b720bb06..c2d891f1 100644 --- a/src/core/sdk/doWhile.ts +++ b/src/core/sdk/doWhile.ts @@ -3,7 +3,8 @@ import { TaskType, DoWhileTaskDef, TaskDefTypes } from "../../common/types"; export const doWhileTask = ( taskRefName: string, terminationCondition: string, - tasks: TaskDefTypes[] + tasks: TaskDefTypes[], + optional?: boolean ): DoWhileTaskDef => ({ name: taskRefName, taskReferenceName: taskRefName, @@ -11,6 +12,7 @@ export const doWhileTask = ( inputParameters: {}, type: TaskType.DO_WHILE, loopOver: tasks, + optional, }); const loopForCondition = (taskRefName: string, valueKey: string) => @@ -19,7 +21,8 @@ const loopForCondition = (taskRefName: string, valueKey: string) => export const newLoopTask = ( taskRefName: string, iterations: number, - tasks: TaskDefTypes[] + tasks: TaskDefTypes[], + optional?: boolean ): DoWhileTaskDef => ({ name: taskRefName, taskReferenceName: taskRefName, @@ -29,4 +32,5 @@ export const newLoopTask = ( }, type: TaskType.DO_WHILE, loopOver: tasks, + optional, }); diff --git a/src/core/sdk/dynamicFork.ts b/src/core/sdk/dynamicFork.ts index 73381fdd..8e440c2d 100644 --- a/src/core/sdk/dynamicFork.ts +++ b/src/core/sdk/dynamicFork.ts @@ -3,7 +3,8 @@ import { TaskType, ForkJoinDynamicDef, TaskDefTypes } from "../../common/types"; export const dynamicForkTask = ( taskReferenceName: string, preForkTasks: TaskDefTypes[] = [], - dynamicTasksInput: string = "" + dynamicTasksInput: string = "", + optional?: boolean ): ForkJoinDynamicDef => ({ name: taskReferenceName, taskReferenceName, @@ -14,4 +15,5 @@ export const dynamicForkTask = ( type: TaskType.FORK_JOIN_DYNAMIC, dynamicForkTasksParam: "dynamicTasks", dynamicForkTasksInputParamName: "dynamicTasksInput", + optional, }); diff --git a/src/core/sdk/event.ts b/src/core/sdk/event.ts index 77aad62f..d0a10493 100644 --- a/src/core/sdk/event.ts +++ b/src/core/sdk/event.ts @@ -3,18 +3,24 @@ import { TaskType, EventTaskDef } from "../../common/types"; export const eventTask = ( taskReferenceName: string, eventPrefix: string, - eventSuffix: string + eventSuffix: string, + optional?: boolean ): EventTaskDef => ({ name: taskReferenceName, taskReferenceName, sink: `${eventPrefix}:${eventSuffix}`, type: TaskType.EVENT, + optional, }); -export const sqsEventTask = (taskReferenceName: string, queueName: string) => - eventTask(taskReferenceName, "sqs", queueName); +export const sqsEventTask = ( + taskReferenceName: string, + queueName: string, + optional?: boolean +) => eventTask(taskReferenceName, "sqs", queueName, optional); export const conductorEventTask = ( taskReferenceName: string, - eventName: string -) => eventTask(taskReferenceName, "conductor", eventName); + eventName: string, + optional?: boolean +) => eventTask(taskReferenceName, "conductor", eventName, optional); diff --git a/src/core/sdk/forkJoin.ts b/src/core/sdk/forkJoin.ts index 7708aa16..23cc1dad 100644 --- a/src/core/sdk/forkJoin.ts +++ b/src/core/sdk/forkJoin.ts @@ -1,4 +1,9 @@ -import { TaskType, ForkJoinTaskDef, TaskDefTypes, JoinTaskDef } from "../../common/types"; +import { + TaskType, + ForkJoinTaskDef, + TaskDefTypes, + JoinTaskDef, +} from "../../common/types"; import { generateJoinTask } from "../generators"; export const forkTask = ( @@ -13,8 +18,9 @@ export const forkTask = ( export const forkTaskJoin = ( taskReferenceName: string, - forkTasks: TaskDefTypes[] + forkTasks: TaskDefTypes[], + optional?: boolean ): [ForkJoinTaskDef, JoinTaskDef] => [ forkTask(taskReferenceName, forkTasks), - generateJoinTask({name:`${taskReferenceName}_join`}), + generateJoinTask({ name: `${taskReferenceName}_join`, optional }), ]; diff --git a/src/core/sdk/http.ts b/src/core/sdk/http.ts index 592f7008..cb0cbde1 100644 --- a/src/core/sdk/http.ts +++ b/src/core/sdk/http.ts @@ -1,13 +1,10 @@ -import { - TaskType, - HttpTaskDef, - HttpInputParameters, -} from "../../common/types"; +import { TaskType, HttpTaskDef, HttpInputParameters } from "../../common/types"; export const httpTask = ( taskReferenceName: string, inputParameters: HttpInputParameters, - asyncComplete = false + asyncComplete?: boolean, + optional?: boolean ): HttpTaskDef => ({ name: taskReferenceName, taskReferenceName, @@ -15,5 +12,6 @@ export const httpTask = ( http_request: inputParameters, }, asyncComplete, + optional, type: TaskType.HTTP, }); diff --git a/src/core/sdk/inline.ts b/src/core/sdk/inline.ts index 387cf21f..5a9f34f2 100644 --- a/src/core/sdk/inline.ts +++ b/src/core/sdk/inline.ts @@ -3,7 +3,8 @@ import { TaskType, InlineTaskDef } from "../../common/types"; export const inlineTask = ( taskReferenceName: string, script: string, - evaluatorType: "javascript" | "graaljs" = "javascript" + evaluatorType: "javascript" | "graaljs" = "javascript", + optional?: boolean ): InlineTaskDef => ({ name: taskReferenceName, taskReferenceName, @@ -12,4 +13,5 @@ export const inlineTask = ( expression: script, }, type: TaskType.INLINE, + optional, }); diff --git a/src/core/sdk/join.ts b/src/core/sdk/join.ts index 18387fce..5cdc038c 100644 --- a/src/core/sdk/join.ts +++ b/src/core/sdk/join.ts @@ -2,10 +2,12 @@ import { TaskType, JoinTaskDef } from "../../common/types"; export const joinTask = ( taskReferenceName: string, - joinOn: string[] + joinOn: string[], + optional?: boolean ): JoinTaskDef => ({ name: taskReferenceName, taskReferenceName, joinOn, type: TaskType.JOIN, + optional, }); diff --git a/src/core/sdk/jsonJq.ts b/src/core/sdk/jsonJq.ts index 4d7ee520..ea805294 100644 --- a/src/core/sdk/jsonJq.ts +++ b/src/core/sdk/jsonJq.ts @@ -2,7 +2,8 @@ import { TaskType, JsonJQTransformTaskDef } from "../../common/types"; export const jsonJqTask = ( taskReferenceName: string, - script: string + script: string, + optional?: boolean ): JsonJQTransformTaskDef => ({ name: taskReferenceName, taskReferenceName, @@ -10,4 +11,5 @@ export const jsonJqTask = ( inputParameters: { queryExpression: script, }, + optional, }); diff --git a/src/core/sdk/kafkaPublish.ts b/src/core/sdk/kafkaPublish.ts index 384ed443..28e3d40a 100644 --- a/src/core/sdk/kafkaPublish.ts +++ b/src/core/sdk/kafkaPublish.ts @@ -6,7 +6,8 @@ import { export const kafkaPublishTask = ( taskReferenceName: string, - kafka_request: KafkaPublishInputParameters + kafka_request: KafkaPublishInputParameters, + optional?: boolean ): KafkaPublishTaskDef => ({ taskReferenceName, name: taskReferenceName, @@ -14,4 +15,5 @@ export const kafkaPublishTask = ( inputParameters: { kafka_request, }, + optional, }); diff --git a/src/core/sdk/setVariable.ts b/src/core/sdk/setVariable.ts index 883ad139..43fc3adb 100644 --- a/src/core/sdk/setVariable.ts +++ b/src/core/sdk/setVariable.ts @@ -2,10 +2,12 @@ import { TaskType, SetVariableTaskDef } from "../../common/types"; export const setVariableTask = ( taskReferenceName: string, - inputParameters: Record + inputParameters: Record, + optional?: boolean ): SetVariableTaskDef => ({ name: taskReferenceName, taskReferenceName, type: TaskType.SET_VARIABLE, inputParameters, + optional, }); diff --git a/src/core/sdk/simple.ts b/src/core/sdk/simple.ts index c488ebcc..3393747b 100644 --- a/src/core/sdk/simple.ts +++ b/src/core/sdk/simple.ts @@ -3,10 +3,12 @@ import { TaskType, SimpleTaskDef } from "../../common/types"; export const simpleTask = ( taskReferenceName: string, name: string, - inputParameters:Record + inputParameters: Record, + optional?: boolean ): SimpleTaskDef => ({ name, taskReferenceName, inputParameters, type: TaskType.SIMPLE, + optional, }); diff --git a/src/core/sdk/subWorkflow.ts b/src/core/sdk/subWorkflow.ts index e10fc7e7..b00755f9 100644 --- a/src/core/sdk/subWorkflow.ts +++ b/src/core/sdk/subWorkflow.ts @@ -3,7 +3,8 @@ import { TaskType, SubWorkflowTaskDef } from "../../common/types"; export const subWorkflowTask = ( taskReferenceName: string, workflowName: string, - version?: number + version?: number, + optional?: boolean ): SubWorkflowTaskDef => ({ name: taskReferenceName, taskReferenceName, @@ -12,4 +13,5 @@ export const subWorkflowTask = ( version, }, type: TaskType.SUB_WORKFLOW, + optional, }); diff --git a/src/core/sdk/switch.ts b/src/core/sdk/switch.ts index 80732da3..d5a023f5 100644 --- a/src/core/sdk/switch.ts +++ b/src/core/sdk/switch.ts @@ -4,7 +4,8 @@ export const switchTask = ( taskReferenceName: string, expression: string, decisionCases: Record = {}, - defaultCase: TaskDefTypes[] = [] + defaultCase: TaskDefTypes[] = [], + optional?: boolean ): SwitchTaskDef => ({ name: taskReferenceName, taskReferenceName, @@ -16,4 +17,5 @@ export const switchTask = ( expression: "switchCaseValue", defaultCase, type: TaskType.SWITCH, + optional, }); diff --git a/src/core/sdk/wait.ts b/src/core/sdk/wait.ts index 64ed3023..da1a0a9e 100644 --- a/src/core/sdk/wait.ts +++ b/src/core/sdk/wait.ts @@ -1,19 +1,29 @@ import { TaskType, WaitTaskDef } from "../../common/types"; -export const waitTaskDuration = (taskReferenceName:string,duration:string):WaitTaskDef =>({ - name:taskReferenceName, - taskReferenceName, - inputParameters:{ - duration - }, - type:TaskType.WAIT +export const waitTaskDuration = ( + taskReferenceName: string, + duration: string, + optional?: boolean +): WaitTaskDef => ({ + name: taskReferenceName, + taskReferenceName, + inputParameters: { + duration, + }, + type: TaskType.WAIT, + optional, }); -export const waitTaskUntil = (taskReferenceName:string,until:string):WaitTaskDef =>({ - name:taskReferenceName, - taskReferenceName, - inputParameters:{ - until - }, - type:TaskType.WAIT -}) \ No newline at end of file +export const waitTaskUntil = ( + taskReferenceName: string, + until: string, + optional?: boolean +): WaitTaskDef => ({ + name: taskReferenceName, + taskReferenceName, + inputParameters: { + until, + }, + type: TaskType.WAIT, + optional, +});